|
This version is still in development and is not considered stable yet. For the latest stable version, please use Korvet 0.19! |
Consumer: Commit Offsets
This page describes how Korvet handles Kafka OffsetCommitRequest operations.
| This page is implementation-oriented and focuses on Kafka-to-Redis mapping details. |
Consumer Commits Offset
Kafka Wire Protocol
OffsetCommitRequest {
group_id: "my-group"
generation_id: 1 // from JoinGroup response
member_id: "consumer-1-uuid"
group_instance_id: null
retention_time_ms: -1 // -1 = use broker default
topics: [
{
name: "orders"
partitions: [
{
partition: 0
committed_offset: 44 // next offset to read
committed_leader_epoch: -1
committed_metadata: "" // accepted but not persisted
}
]
}
]
}
OffsetCommitResponse {
throttle_time_ms: 0
topics: [
{
name: "orders"
partitions: [
{
partition: 0
error_code: 0 // NONE
}
]
}
]
}
Redis Commands Detail
Offset commits have two storage effects in the current implementation:
-
the broker acknowledges previously delivered Redis Stream entries with
XACK -
the broker persists the committed Kafka offset and commit time in the committed-offset store keyed by
(groupId, streamKey)
The Redis store keeps canonical offset values as bare numbers so older brokers can
continue to read them during a rolling upgrade. A parallel per-group timestamp hash
stores each value as <offset>|<commit-epoch-millis>. A generated hash tag co-locates
both hashes in Redis Cluster, allowing each offset/timestamp pair to be written
atomically. Timestamp metadata is non-authoritative: a malformed or unavailable
metadata hash never fails a canonical offset write or read.
The timestamp represents the latest commit observed by a timestamp-aware broker.
Commits handled exclusively by an older broker during a rolling upgrade cannot
advance or invalidate that observation, including a legacy re-commit of the same
offset. Untimestamped writes through the current store clear prior metadata, as do
offset deletion and migration replacement.
Kafka commit metadata is accepted on the request but is not persisted, so it cannot
be returned by a later OffsetFetch.
Implementation Notes
-
Offset Semantics: Kafka commits the next offset to read. The broker acknowledges every tracked delivered Redis message whose encoded offset is strictly less than the committed offset.
-
Tracked Delivery State: Delivered Redis message IDs are tracked in broker memory per
(groupId, memberId, topic, partition)so commit can batch acknowledgments efficiently. -
Committed Offset Persistence:
OffsetFetchreads the committed-offset store first, so successful commits are not represented solely by Redis consumer-group metadata. The Admin API also reads the persisted commit time for itslastCommitTimestampfield. -
Batch / Retry Behavior: The handler first tries batched acknowledgments and falls back to per-partition retries if needed.