FileOffsetStore
Generated page
Model gemma-mtp, commit f508a4b65b3f, 2026-08-16, sources: 6. Edit the code or the hand-written documentation instead.
What this module is responsible for
The FileOffsetStore provides a simple implementation of the OffsetStore interface, designed for consumers that need to persist their position in a log on a local volume. It ensures that a consumer's position outlives the process by storing the Offset in a file, allowing for recovery after a restart FileOffsetStore.kt:27-29.
Diagram
FileOffsetStore
The FileOffsetStore is a concrete implementation of OffsetStore that uses the local file system to track progress. It maps a combination of TopicName and PartitionId to a specific file path within a configured directory FileOffsetStore.kt:61-64.
Atomic move via temporary files
To prevent data corruption, the save operation does not overwrite the existing offset file directly. Instead, it writes the new offset to a sibling file with a .tmp extension and then performs an atomic move using StandardCopyOption.ATOMIC_MOVE FileOffsetStore.kt:55-57. This ensures that a partially written file (e.g., a truncated "12" becoming "1") does not result in a valid but incorrect offset that could cause silent record replays FileOffsetStore.kt:23-25.
At-least-once checkpointing
The system guarantees at-least-once delivery by performing checkpointing only after the collector has successfully handled a batch of records README.md:35-36. If a failure occurs between the processing of the batch and the call to save, the consumer will restart from the previous offset and replay the batch README.md:36.
File-based position recovery
When a consumer restarts, it uses the load function to retrieve the last known Offset from the file system FileOffsetStore.kt:34-46. Because the broker does not store consumer positions, the responsibility for maintaining this state lies entirely with the client/consumer feature-subscribe-and-publish.md:30-31.
Key files
| File | Lines | What is there |
|---|---|---|
…/common/FileOffsetStore.kt | 27-64 | Implementation of FileOffsetStore including load, save, and file path resolution. |
Behaviour that surprises
- The
savemethod uses a temporary file andStandardCopyOption.ATOMIC_MOVEto ensure that a position saved halfway through a write does not result in a truncated offsetFileOffsetStore.kt:55-57. checkpointingis designed for at-least-once semantics; saving the offset before handling the batch would result in at-most-once delivery, potentially skipping recordsREADME.md:35-36.- The
loadfunction returnsnullif the offset file does not exist, which is treated as a valid starting point for a new consumerFileOffsetStore.kt:40.