Explain ProcessorContext.schedule() and the difference between PunctuationType.STREAM_TIME and WALL_CLOCK_TIME punctuators.
answer
- schedule(Duration, PunctuationType, Punctuator) -> Cancellable
- STREAM_TIME = max record timestamp, data-driven
- WALL_CLOCK_TIME = system clock, fires on idle too
- stream-time = deterministic/reprocessable; wall-clock = not
- wall-clock cadence bounded by poll loop, not exact
- punctuate runs on stream thread, can forward + use stores
basics
~20 sschedule() registers a punctuator — a callback that fires periodically. STREAM_TIME advances by the timestamps of records flowing through, so it only fires when data arrives and progresses. WALL_CLOCK_TIME advances by the system clock, so it fires on a real-time interval even if no records arrive.
solid answer
~50 scontext.schedule(Duration interval, PunctuationType type, Punctuator p) registers a periodic callback returning a Cancellable. The punctuator's punctuate(long timestamp) runs on the stream thread, between record processing, so it can safely use state stores and call context.forward(). The two clocks differ fundamentally: STREAM_TIME is driven by record timestamps (event time / the max observed timestamp) — it only advances and fires as records arrive, and it pauses when the input is idle, making it deterministic and reprocessing-friendly. WALL_CLOCK_TIME is driven by the system clock (the producer/system time), firing every interval of real time regardless of input — useful for flushing, timeouts, heartbeats, or emitting even when data stops. WALL_CLOCK punctuators only fire when the thread polls/loops, so the actual cadence is bounded by poll frequency and isn't guaranteed exact. Cancel via the returned Cancellable, e.g., in close(). On reprocessing, stream-time replays identically; wall-clock does not.
go deeper
Know punctuators are scheduled callbacks and that the two types use record time vs. real time.
Explain schedule()'s signature, Cancellable, and that punctuate runs on the stream thread and can forward.
Contrast stream-time (data-driven, deterministic, idle-pausing) vs. wall-clock (real-time, fires on idle, best-effort cadence) and pick correctly per use case.
Reason about reprocessing determinism, idle-partition/stream-time advancement, max.poll.interval interactions, and idempotent catch-up firing.
**Punctuation** is how the Processor API runs **time-driven logic** — code that should execute on a schedule rather than per input record. You register it in `init()`: ``` this.cancellable = context.schedule( Duration.ofSeconds(10), PunctuationType.STREAM_TIME, // or WALL_CLOCK_TIME timestamp -> { /* punctuate body */ }); ``` `schedule(Duration interval, PunctuationType type, Punctuator punctuator)` returns a **`Cancellable`**. The **Punctuator** is a functional interface with one method, `void punctuate(long timestamp)`, where `timestamp` is the time at which the punctuation fired (in the chosen time domain). **Execution model.** A punctuator runs on the **same single stream thread** that processes records for the task, **between** `process()` calls — never concurrently. Because of that, it may safely read/write the task's **state stores** and call **`context.forward(...)`** to emit results. Typical uses: emit an aggregation buffered in a store, expire/clean up old keys, send a heartbeat, or flush batched output. **The two time domains — the core distinction:** **1. `PunctuationType.STREAM_TIME` (event/stream time).** "Stream time" is the **maximum record timestamp** the task has observed so far (record timestamps come from the `TimestampExtractor`, default = the record's embedded event timestamp). A stream-time punctuator fires when stream time has **advanced** by at least `interval` since the last firing. Consequences: - It **only advances when records arrive** — if the input is idle, stream time is frozen and the punctuator **does not fire**. - It can fire **multiple times in one go** if a record jumps stream time forward by several intervals (it catches up). - It is **deterministic and reprocessing-safe**: replaying the same input produces the same punctuation schedule, because it depends only on data, not the wall clock. Ideal for windowed/event-time emission and tests. **2. `PunctuationType.WALL_CLOCK_TIME` (system/processing time).** Driven by the **system wall-clock** (`Time` abstraction, effectively `System.currentTimeMillis()`). It fires every `interval` of **real elapsed time**, **regardless of whether records arrive**. Consequences: - Works even when the input topic is **idle** — good for timeouts, periodic flushes, liveness/heartbeats, or emitting "no data" alerts. - The cadence is **best-effort, not exact**: a wall-clock punctuator can only fire when the stream thread checks the clock, which happens around its **poll loop** iterations. A long `process()`, a long poll, or `max.poll.interval` blocking delays it. So treat the interval as a lower bound, not a precise timer. - It is **non-deterministic** under reprocessing — replay timing depends on how fast the data is consumed, so the same input yields different wall-clock punctuation. Avoid it where determinism matters. **Cancellation & lifecycle.** Keep the returned `Cancellable` and call `cancel()` when appropriate (e.g., in the processor's `close()`); punctuators are also torn down when the task closes/migrates. Re-scheduling on each `init()` is normal because `init()` runs per task start. **Common patterns & edge cases:** - *Aggregate-then-emit:* buffer per-key state in a store, and on STREAM_TIME punctuation emit/flush completed windows — classic for custom windowing. - *Idle timeout / sessionization:* WALL_CLOCK_TIME to close sessions when no new events arrive (stream time wouldn't advance to trigger it). - *Don't block in punctuate:* it runs on the processing thread; slow work stalls record processing and can trip `max.poll.interval.ms`. - *Stream-time starvation:* if some partitions are idle, stream time may not advance as expected; the `task.idle.ms`/idling config and the `TimestampExtractor` influence this. - *Multiple firings:* design punctuate to be idempotent / loop-safe since stream-time can catch up several intervals at once. - *Forwarding timestamp:* records forwarded from a punctuator get a timestamp you set (often the punctuation `timestamp`), which matters for downstream time semantics.
- Why might a STREAM_TIME punctuator never fire even though your app is running?Stream time only advances when records with newer timestamps arrive. If the input is idle (or partitions are idle and time can't advance), stream time is frozen, so the punctuator won't fire. Use WALL_CLOCK_TIME if you need it to fire on idle input.
- Is the wall-clock interval exact?No. A WALL_CLOCK_TIME punctuator can only fire when the stream thread checks the clock around its poll loop, so a long process(), long poll, or blocking delays it. Treat the interval as a lower bound.
- Which punctuation type should you use for reprocessing/testing determinism, and why?STREAM_TIME — it depends only on record timestamps, so replaying the same input yields the same punctuation schedule. WALL_CLOCK_TIME depends on real consumption speed and is non-deterministic.
saying these in an interview costs you the question
- Claiming WALL_CLOCK_TIME fires on a precise timer thread — it fires around poll-loop iterations and is best-effort.
- Saying STREAM_TIME fires on a fixed real-time schedule — it is driven by record timestamps and pauses on idle input.
- Believing punctuators run on a separate thread — they run on the stream thread, between records.
- Thinking stream-time punctuation never catches up — one record can trigger several missed intervals at once.
- Using STREAM_TIME for idle-timeout logic — it won't fire when data stops.