skip to content

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

level: middleimportance: must knowfreq 60%

answer

  1. interface extract(record, partitionTime)
  2. Fail = throw, LogAndSkip = warn+drop, UsePartitionTime = inherit
  3. Wallclock = System.currentTimeMillis()
  4. default.timestamp.extractor / Consumed.withTimestampExtractor
  5. custom extractor reads payload field

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().

solid answer

~50 s

TimestampExtractor is the SPI Kafka Streams calls for every record to determine its timestamp, which then drives windowing, joins, and stream time. Built-in implementations: ExtractRecordMetadataTimestamp is the base that returns the record's embedded message timestamp. FailOnInvalidTimestamp (the default) extends it and throws a StreamsException if the timestamp is negative (e.g., from a message format predating timestamps, or a buggy producer), failing the application. LogAndSkipOnInvalidTimestamp returns the invalid timestamp as-is but logs a warning — a negative timestamp causes the record to be effectively dropped from time-based operations. WallclockTimestampExtractor ignores the embedded timestamp entirely and returns System.currentTimeMillis(), giving processing-time semantics. You set it via default.timestamp.extractor in StreamsConfig, or per-Consumed via Consumed.with(...). A common pattern is a custom extractor that pulls event time from the payload, falling back to the previous record's timestamp (the second argument) when the field is missing.

go deeper

for a junior

Know there's a TimestampExtractor and that the default reads the record's timestamp; Wallclock uses current time.

for a middle

Distinguish Fail vs LogAndSkip vs UsePartitionTime vs Wallclock and how to configure them.

for a senior

Explain the extract() signature, partitionTime fallback, and writing a payload-based custom extractor.

for a principal

Decide extractor strategy per data quality profile and reason about how a custom/wall-clock extractor changes determinism and window correctness across the pipeline.

## What it is `org.apache.kafka.streams.processor.TimestampExtractor` is a single-method interface: ``` long extract(ConsumerRecord<Object,Object> record, long partitionTime) ``` Kafka Streams invokes it for **every** record it reads to decide what timestamp that record carries *inside* the topology. That timestamp drives windowed aggregations, windowed joins, suppression, and the advancement of **stream time**. `partitionTime` (formerly `previousTimestamp`) is the highest timestamp seen so far on that partition — useful as a fallback. ## Built-in extractors All except the wall-clock one extend `ExtractRecordMetadataTimestamp`, which returns the record's metadata timestamp (`ConsumerRecord.timestamp()` — i.e. event time under `CreateTime`, ingestion time under `LogAppendTime`). They differ only in how they handle an **invalid (negative)** timestamp: - **FailOnInvalidTimestamp** (default): throws `StreamsException` on a negative timestamp. Records written by very old producers (message format v0) or buggy clients can carry `-1`. This fails fast so you notice bad data. - **LogAndSkipOnInvalidTimestamp**: logs a warning and returns the negative value. Streams treats records with negative timestamps as not processable for time-based ops, so they are effectively skipped rather than crashing the app — more tolerant of dirty data. - **UsePartitionTimeOnInvalidTimestamp**: on an invalid timestamp, returns `partitionTime` (the current partition/stream time) instead of failing or skipping — keeps the record but inherits the latest known time. - **WallclockTimestampExtractor**: ignores the record completely and returns `System.currentTimeMillis()`. This gives **processing-time** semantics — windows bucket by when the stream thread sees the record, not when the event happened. ## Configuring it - Globally: `default.timestamp.extractor` in `StreamsConfig` (e.g. `WallclockTimestampExtractor.class`). - Per source: `Consumed.with(serde, serde).withTimestampExtractor(new MyExtractor())`. ## Custom extractors The most common real-world need is extracting event time from the payload (the broker timestamp may be ingestion/skewed). You implement `extract` to parse the value, and on a missing/invalid field return `partitionTime` so the record stays close to stream time instead of poisoning windows. ## Edge cases - Returning a negative number from a custom extractor triggers the same invalid-timestamp handling. - Wall-clock extractor breaks reproducibility — replays bucket old data into 'now'. - A future-dated timestamp advances stream time far ahead, which can prematurely close earlier windows and drop subsequent in-order records as late.

  • What is the second argument to extract() and how is it useful?
    partitionTime (previously previousTimestamp) — the highest timestamp seen so far on that partition. Custom extractors return it as a fallback when the payload field is missing or invalid, keeping the record aligned with current stream time instead of failing.
  • When would you choose LogAndSkip over the default FailOnInvalidTimestamp?
    When the topic may contain records from legacy producers or buggy clients with negative timestamps and you'd rather tolerate/skip them than crash the whole application.

saying these in an interview costs you the question

  • Saying WallclockTimestampExtractor reads the record timestamp (it ignores it entirely).
  • Claiming LogAndSkip throws an exception (it logs and continues).
  • Forgetting the default is FailOnInvalidTimestamp.
  • Thinking 'invalid' means any odd value rather than specifically a negative timestamp.

context