skip to content

Time Semantics

Event time versus processing time, timestamp extractors, and how stream time advances with out-of-order records. Interviewers ask because the choice of time silently changes every windowed result.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What are the three notions of time in Kafka Streams (event time, processing time, ingestion time), and how do they differ?

level: juniorimportance: must knowfreq 70%

answer

  1. event=happened, ingestion=broker append, processing=wall-clock
  2. CreateTime vs LogAppendTime topic config
  3. Streams defaults to event time
  4. event time = deterministic replay
  5. processing time = WallclockTimestampExtractor

basics

~20 s

Event time = when the event actually happened (set by the producer). Ingestion time = when the broker appended the record. Processing time = when Kafka Streams processes it. They differ because of network and processing delays.

solid answer

~40 s

Kafka Streams distinguishes three times. Event time is when the event occurred at the source, embedded in the record (typically the message timestamp set by the producer or a field in the payload). Ingestion time is when the Kafka broker wrote the record to the partition; it is set by the broker when the topic's message.timestamp.type is LogAppendTime. Processing time is wall-clock time on the stream thread when the operator handles the record. Most Streams operations (windowing, joins) default to event time, which gives deterministic, reproducible results regardless of when or how fast you reprocess. Processing time is non-deterministic and depends on lag and machine clocks. Which one you get is governed by the configured TimestampExtractor and, for ingestion time, by the broker topic config message.timestamp.type.

go deeper

for a junior

Know the three definitions and that event time = when it happened, processing time = when it's handled.

for a middle

Connect each time to its source: producer timestamp, broker append, stream thread clock; know message.timestamp.type.

for a senior

Explain why event time gives deterministic replay and how the TimestampExtractor + topic config jointly decide the effective time.

for a principal

Reason about choosing CreateTime vs LogAppendTime per topic for downstream event-time correctness across a pipeline, and the irreversibility of LogAppendTime.

## The problem A record flows through several stages: it is created at a source, sent to a Kafka broker, stored in a partition, then later read and processed by a Kafka Streams application. There are several distinct 'times' we could associate with the record, and conflating them causes subtle bugs in windowing and joins. ## The three times - **Event time**: the moment the event actually occurred in the real world. Example: a user clicked 'buy' at 12:00:00.000. This is usually carried in the Kafka record's timestamp field (set by the producer) or inside the message payload (e.g., a JSON `createdAt` field). Event time is the source of truth for *what happened when*. - **Ingestion time** (a.k.a. append/broker time): the moment the **broker** appended the record to its log. Controlled by the topic config `message.timestamp.type` — when set to `LogAppendTime`, the broker overwrites the producer-supplied timestamp with its own clock at append. When set to `CreateTime` (the default), the producer's timestamp is kept (which usually approximates event time). - **Processing time**: the wall-clock time on the Streams instance when an operator processes the record. This depends on consumer lag, restarts, and reprocessing, so it is **non-deterministic**. ## Why it matters Kafka Streams uses *event time* by default for time-driven operations (windowed aggregations, windowed joins, suppression). Event time makes results **deterministic and reproducible**: re-running the same input always yields the same windows, regardless of when you run it or how far behind you are. If you instead keyed everything off processing time, a backlog replay would bucket old events into 'now' windows, producing different results each run. ## How Streams gets the time The `TimestampExtractor` interface decides which time becomes the record's timestamp inside Streams. The default, `FailOnInvalidTimestamp`, returns the record's embedded message timestamp (event time if `CreateTime`, ingestion time if `LogAppendTime`). To force processing time you use `WallclockTimestampExtractor`. To pull a custom field from the payload you write your own extractor. ## Edge cases - A topic with `LogAppendTime` cannot give you true event time downstream — the original producer timestamp is gone. Choose `CreateTime` if you need event-time semantics. - Producer timestamps can be wrong (clock skew, bad clients), leading to negative or future timestamps; that interacts with the extractor's error handling.

  • Which topic config decides whether the record timestamp is producer-set or broker-set?
    message.timestamp.type — CreateTime (default) keeps the producer timestamp (~event time); LogAppendTime overwrites it with the broker's append time (ingestion time).
  • Why does Kafka Streams prefer event time for windowing?
    It makes results deterministic and reproducible: the same input always produces the same windows regardless of when or how fast you reprocess, unlike processing time which depends on lag and machine clocks.

saying these in an interview costs you the question

  • Saying processing time is the default for windowing (event time is the default).
  • Claiming ingestion time and event time are the same thing.
  • Thinking the producer can never affect the record timestamp (CreateTime keeps it).
  • Assuming you can recover event time from a LogAppendTime topic.

context

open as a page

What is the TimestampExtractor interface in Kafka Streams, and what do FailOnInvalidTimestamp, LogAndSkipOnInvalidTimestamp, and WallclockTimestampExtractor do?

level: middleimportance: must knowfreq 60%

basics

~10 s

TimestampExtractor tells Streams which timestamp to use per record. FailOnInvalidTimestamp (default) reads the record timestamp and throws on a negative one; LogAndSkipOnInvalidTimestamp logs and skips bad ones; WallclockTimestampExtractor ignores the record and uses System.currentTimeMillis().

open as a page

What is 'stream time' in Kafka Streams, how does it advance, and why does it never go backwards?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Stream time is the app's notion of the current event-time clock per task: the maximum record timestamp seen so far. It only moves forward — a later record with an older timestamp does not pull it back. It drives window closing and lateness decisions.

open as a page

How does Kafka Streams handle out-of-order and late records in windowed operations, and what role does the grace period play?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Out-of-order records (timestamp below stream time but within the grace period) are still added to their window and update the result. Once stream time passes window end plus the grace period, the window closes and later-arriving records for it are dropped as 'late'.

open as a page

What problem does max.task.idle.ms solve in Kafka Streams, and what are the trade-offs of tuning it?

level: principalimportance: should knowfreq 35%

basics

~10 s

When a task reads several partitions, max.task.idle.ms tells Streams how long to pause/wait for an empty-but-active partition to deliver data before processing ahead with other partitions. It trades join/merge time-ordering correctness against latency.

open as a page