What are the three notions of time in Kafka Streams (event time, processing time, ingestion time), and how do they differ?
answer
- event=happened, ingestion=broker append, processing=wall-clock
- CreateTime vs LogAppendTime topic config
- Streams defaults to event time
- event time = deterministic replay
- processing time = WallclockTimestampExtractor
basics
~20 sEvent 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 sKafka 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
Know the three definitions and that event time = when it happened, processing time = when it's handled.
Connect each time to its source: producer timestamp, broker append, stream thread clock; know message.timestamp.type.
Explain why event time gives deterministic replay and how the TimestampExtractor + topic config jointly decide the effective time.
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.