skip to content

In Flink, what is the difference between event time and processing time?

level: juniorimportance: must knowfreq 82%

answer

  1. two clocks disagree
  2. one travels with the record
  3. the other lives on the TaskManager
  4. replay gives the same windows
  5. no global time switch in 2.x

basics

~20 s

Event time is the timestamp carried inside the record, set when the event happened. Processing time is the clock on the machine running the operator. Event time gives reproducible, order-independent results; processing time gives lower latency and no correctness guarantees.

solid answer

~50 s

**Event time** is a timestamp that travels with the record — the moment the click, payment or sensor reading actually occurred. Flink extracts it with a `TimestampAssigner` and uses it to place records into windows and to fire event-time timers. **Processing time** is simply `System.currentTimeMillis()` on the TaskManager that happens to be executing the operator. The practical difference is reproducibility. Replay a Kafka topic through an event-time job and every window contains exactly the same records as the original live run, no matter how fast or out of order the replay is. The same replay under processing time bunches everything into a handful of windows, because the only clock that matters is the wall clock at the moment of processing. Processing time is cheaper and lower latency: no watermarks, no waiting for stragglers, no late data. Use it for coarse operational metrics or heartbeats where a slightly wrong bucket costs nothing. Use event time for anything a business will reconcile.

code

java · 5 lines
java
WatermarkStrategy<Order> strategy = WatermarkStrategy
    .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20))
    .withTimestampAssigner((order, recordTs) -> order.getPlacedAtMillis());

DataStream<Order> orders = env.fromSource(kafkaSource, strategy, "orders");

go deeper

for a junior

Be ready to state, in one sentence each, that event time comes from the record and processing time comes from the machine clock, and to say which one survives a replay unchanged.

for a middle

Explain how Flink actually obtains the event timestamp — a WatermarkStrategy with a TimestampAssigner — and why event time forces the engine to introduce watermarks while processing time does not.

for a senior

Show judgment about which pipelines deserve event time. Be able to argue the latency and state cost of waiting for stragglers, and to name workloads where processing time is genuinely the correct engineering choice.

for a principal

Own the consequence across the platform: if backfills and live streams must reconcile, event time is not optional, and the whole ingestion chain — producers stamping records, connectors preserving timestamps, schema conventions — has to guarantee a trustworthy timestamp field.

## Three clocks, not one A streaming record can be stamped by three different clocks, and the whole difficulty of stream processing comes from them disagreeing. **Event time** is the time the event happened in the real world, recorded in the record itself — a field in the JSON payload, an Avro column, or the Kafka record's own timestamp. It is fixed forever once the record is created. A mobile phone that was offline in a tunnel emits a purchase with event time 09:00 even if the record reaches the broker at 11:30. **Ingestion time** is the time the record entered the Flink pipeline — assigned at the source operator. It is monotonically increasing per source and immune to out-of-order arrival, but it is not the truth about the event, and it is not reproducible on replay either. **Processing time** is the local system clock of whichever TaskManager executes the operator at the instant it touches the record. Different subtasks on different machines will read slightly different values, and a restart changes them entirely. ## What event time buys you Determinism. Under event time, a five-minute tumbling window over `[09:00, 09:05)` contains exactly the records whose own timestamps fall in that range — regardless of whether they arrived in order, arrived hours late, or were replayed from the beginning of the topic during a backfill. That is what makes it possible to: recompute yesterday's aggregates and get the same numbers; run a batch backfill and a live stream over the same logic and reconcile them; and reason about correctness at all when the network reorders records. This is also why event time is the only sensible choice for anything financial, billing-related, or user-facing. If a customer's session window depends on when your cluster happened to schedule the operator, the result is not a fact about the customer. ## What event time costs you Somebody has to decide when a window is complete. Under processing time the answer is trivial — when the clock passes the window end. Under event time the engine has no idea whether a straggler is still in flight, so Flink introduces **watermarks**: a watermark of value `T` flowing through the stream is the pipeline's assertion that no further record with a timestamp at or before `T` is expected. Windows fire when the watermark passes their end. Everything else in event-time processing — bounded out-of-orderness, allowed lateness, late side outputs, idle sources — exists to manage that assertion. So event time buys correctness and pays in latency (you wait for the watermark), state (open windows are held until they fire), and complexity (you have to configure how long to wait). ## How Flink learns the timestamp You attach a `WatermarkStrategy` that carries a `TimestampAssigner`: ```java WatermarkStrategy<Order> strategy = WatermarkStrategy .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((order, recordTs) -> order.getPlacedAtMillis()); ``` The second argument `recordTs` is the timestamp the source already attached — for the Kafka connector, the Kafka record's own timestamp (`ConsumerRecord.timestamp()`: the producer's create time or the broker's append time, depending on how the topic is configured). If your payload has no usable time field, that source-provided value is your fallback. It is stored with the record, so it survives a replay, but it says when the record was produced or appended, not necessarily when the event happened. Attach the strategy in `env.fromSource(source, strategy, "kafka")` where possible rather than calling `assignTimestampsAndWatermarks` afterwards: the source generates watermarks per split, which behaves far better when a subtask reads several Kafka partitions. ## Processing time is still a legitimate choice Do not treat processing time as the beginner's option. It is the right tool when the question you are answering is literally about wall-clock behaviour: "how many errors did we log in the last minute", alerting heartbeats, rate limiting, timeouts on inactivity. It has no watermarks, no late data, no lateness dial, and its latency is as low as the engine can go. `ProcessFunction` lets you register processing-time timers alongside event-time ones, so a single job can use both — a common pattern is an event-time computation with a processing-time timeout to flush a stuck key. ## Version note This answer assumes Flink 2.3. Older Flink selected the semantics globally with `env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)`. Event time became the default in 1.12, when that setter was deprecated, and Flink 2.0 removed both the method and `TimeCharacteristic`. Ingestion time is no longer a separate characteristic either: you get it by stamping records with the wall clock right after the source, which is exactly what Flink's `IngestionTimeAssigner` does. In Flink 2.3 the choice is made per operation: attach a `WatermarkStrategy` with a timestamp assigner and use event-time windows and timers, or use processing-time windows and timers and attach nothing.

  • If event time is more correct, why would you ever pick processing time?
    Because it is free. Processing time needs no watermarks, so windows fire the instant the clock passes their end — no waiting for stragglers, no late data, no state held open. For operational metrics, heartbeats, inactivity timeouts and rate limiting, the question genuinely is about wall-clock behaviour, so the extra machinery buys nothing.
  • Where does ingestion time fit in modern Flink?
    It is no longer a separate time characteristic. You get it by stamping each record with the wall clock as it leaves the source: Flink ships `IngestionTimeAssigner` for this, plugged in with `withTimestampAssigner(ctx -> new IngestionTimeAssigner<>())`, usually on `forMonotonousTimestamps()` because the stamps only increase. It removes out-of-orderness but is not reproducible across replays.
  • Can one Flink job mix both?
    Yes. A `KeyedProcessFunction` can register event-time timers via `ctx.timerService().registerEventTimeTimer(...)` and processing-time timers via `registerProcessingTimeTimer(...)` on the same key. A common pattern is an event-time aggregation with a processing-time safety valve that flushes a key whose watermark has stopped advancing.

Event time is the postmark on a letter; processing time is the moment you happened to open it. Sorting your mail by postmark gives the same answer no matter when you get to the pile.

saying these in an interview costs you the question

  • Says event time is the time the record reached Flink
  • Claims processing time gives the same results on replay
  • Thinks Flink reads the timestamp automatically without an assigner
  • Says you must call setStreamTimeCharacteristic in current Flink versions
  • Treats ingestion time as identical to event time

context