How does the Confluent S3 sink connector lay out files in object storage, and how can it achieve exactly-once delivery to S3?
answer
- flush.size or rotate.interval
- TimeBasedPartitioner → year=/month=/day=/hour=
- filename encodes starting offset
- atomic S3 PUT overwrite = exactly-once
- no Wallclock extractor for EOS
basics
~20 sThe S3 sink groups records by partitioner (by Kafka partition or by time) and writes batched objects (Parquet/Avro/JSON) when flush.size or a rotate interval is hit. It achieves exactly-once by naming files deterministically from offsets, so replays overwrite the same object instead of duplicating.
solid answer
~40 sThe S3 sink buffers records per topic-partition and emits an object when `flush.size` records accumulate or `rotate.interval.ms`/`rotate.schedule.interval.ms` elapses. A **partitioner** controls the key prefix: `DefaultPartitioner` uses `topic/partition`, while `TimeBasedPartitioner` derives `year=/month=/day=/hour=` paths from a record/wallclock/record-field timestamp — giving Hive-style partitions analytics engines can prune. Formats include Parquet, Avro, JSON, and raw bytes via the `format.class`. **Exactly-once** to S3 is achieved because each output file is named deterministically from the starting Kafka **offset** of its records; on failure/replay the connector rewrites the *same* object key, and S3's atomic single-object PUT means the replay overwrites rather than appends. The connector only advances committed offsets after the upload succeeds, so the file set is consistent with offsets. This requires deterministic partitioning (no wall-clock partitioner) for the guarantee to hold.
go deeper
Know the sink batches records into files in S3 with time-based folders and supports Parquet.
Explain flush.size/rotation, partitioner types, and the offset-in-filename exactly-once mechanism.
Detail why deterministic partitioning is required and trade-offs of file size vs. small-file problems.
Weigh S3 sink vs. native Iceberg sink, compaction strategy, and lake table-format governance at scale.
## The challenge Object stores like S3 have no transactions and no append — you PUT whole objects. So a sink that wants *exactly-once* output cannot rely on a database transaction. The Confluent S3 sink solves this with **deterministic, offset-based file naming**. ## File layout The sink buffers incoming records **per topic-partition** in memory and flushes an object when either: - `flush.size` records have accumulated, or - a rotation timer fires: `rotate.interval.ms` (based on record timestamps) or `rotate.schedule.interval.ms` (wall-clock). The **partitioner** decides the S3 key prefix: - `DefaultPartitioner` → `topics/<topic>/partition=<p>/...` - `FieldPartitioner` → groups by a record field value (e.g., `region=eu`). - `TimeBasedPartitioner` → Hive-style `year=2026/month=06/day=30/hour=14/...` using `timestamp.extractor` = `Record` (event time), `RecordField`, or `Wallclock`. These directory-style partitions let query engines (Athena, Trino, Spark) **prune** scans by date. Output **format** is set by `format.class`: `ParquetFormat` (columnar, best for analytics), `AvroFormat`, `JsonFormat`, or `ByteArrayFormat`. ## How exactly-once works 1. Each flushed file's name **encodes the Kafka offset** of the first record it contains, e.g., `...+0000000000.parquet` where the number is the starting offset. 2. On a crash, uncommitted offsets cause Connect to **replay** the same records. Because the partitioner is deterministic, those same records map to the **same S3 object key**. 3. S3 PUT of a single object is **atomic** — the replay simply **overwrites** the previous (possibly partial) object with the complete one. No duplicate file, no partial append. 4. The sink commits offsets only after the upload, keeping the S3 file set consistent with committed offsets. ## Required conditions - **Deterministic partitioning**: a `Wallclock` timestamp extractor breaks the guarantee because the same record could land in different time partitions across replays, producing different object keys and thus duplicates. Use record/field-based extractors for exactly-once. - The guarantee is **per the connector's own writes**; downstream readers must still read complete files (the sink writes atomically so partial files aren't exposed). ## Why analysts care Hive-style time partitioning plus Parquet means cheap, prunable, columnar scans; offset-deterministic naming means re-running after an incident won't double-count revenue. It's the practical bridge from streaming to a query-able lake.
- Why does using a Wallclock timestamp extractor break exactly-once for the S3 sink?Because the output object key depends on when the record is processed, not on its content. On replay the same record can land in a different time partition, producing a new file key instead of overwriting the original — so you get duplicates rather than idempotent rewrites.
- Why is Parquet preferred over JSON as the S3 sink output format for analytics?Parquet is columnar and compressed, so engines read only the columns they need and prune row groups via statistics, drastically cutting scan cost and bytes read versus row-oriented JSON.
saying these in an interview costs you the question
- Claiming S3 supports appends or transactions — it does not; the sink relies on atomic full-object PUT.
- Saying exactly-once works with any partitioner including Wallclock.
- Thinking the sink writes one file per record — it batches by flush.size/rotation.