What is the TimestampExtractor interface in Kafka Streams, and what do FailOnInvalidTimestamp, LogAndSkipOnInvalidTimestamp, and WallclockTimestampExtractor do?
answer
- interface extract(record, partitionTime)
- Fail = throw, LogAndSkip = warn+drop, UsePartitionTime = inherit
- Wallclock = System.currentTimeMillis()
- default.timestamp.extractor / Consumed.withTimestampExtractor
- custom extractor reads payload field
basics
~10 sTimestampExtractor 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 sTimestampExtractor 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
Know there's a TimestampExtractor and that the default reads the record's timestamp; Wallclock uses current time.
Distinguish Fail vs LogAndSkip vs UsePartitionTime vs Wallclock and how to configure them.
Explain the extract() signature, partitionTime fallback, and writing a payload-based custom extractor.
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.