Skip to content

SS-455 - source: commit Kafka offsets on change and on a timer - #39462

Open
martykulma wants to merge 2 commits into
maz-offset-committed-per-exportfrom
maz-kafka-offset-commit-dedupe
Open

martykulma wants to merge 2 commits into
maz-offset-committed-per-exportfrom
maz-kafka-offset-commit-dedupe

Conversation

@martykulma

@martykulma martykulma commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Each worker commits only when its offsets change, plus on a timer. The timer is needed because the consumer only ever assigns partitions, which results in the brokers treating the group as standalone. For standalone groups, brokers expire a partition's offset offsets.retention.minutes after its last commit, whether or not it is still current. The interval is the kafka_offset_commit_refresh_interval dyncfg: 10 minutes by default, 10 seconds in CI, zero disables it. A failed commit is retried on the next resume upper.

🤖 Generated with Claude Code
👨 Improved by human

@martykulma
martykulma added this pull request to stack #39464 October 1, 2026 18:35
@martykulma martykulma changed the title source: Kafka only comits offsets for a partition if they've changed SS-455 - source: Kafka only commits offsets for a partition if they've changed Oct 1, 2026
@linear-code

linear-code Bot commented Oct 1, 2026

Copy link
Copy Markdown

SS-455

@martykulma
martykulma marked this pull request as ready for review October 1, 2026 19:30
@martykulma
martykulma requested a review from a team as a code owner October 1, 2026 19:30
@def-

def- commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

QA LLM Review

1. MEDIUM -- Idle partitions are never re-committed, so the broker expires their consumer-group offsets

src/storage/src/source/kafka.rs:1163

With the offsets != self.committed_offsets check, a worker whose partitions receive no new messages commits them once and never again, even while the rest of the source keeps advancing. Materialize consumes with assign(), which makes it a standalone consumer group, and brokers delete a standalone group's offset for a partition offsets.retention.minutes (7 days by default) after that partition's last commit. Lag-monitoring tools, which the Kafka source docs say these commits exist for, then lose those partitions.

Details

On main, reclock_committed_upper calls tx.send_replace every time it runs (source_reader_pipeline.rs:685 at the merge base), so every worker re-commits its offsets about once per tick, and that refreshes the commit timestamp. On #39461, any change to ResumeUppers still re-commits every worker's partitions. This PR takes away that last refresh for idle partitions. Two examples: a keyed topic where some partitions get no writes for a week, or a cluster with at least as many workers as partitions, where each worker owns at most one partition. After the retention window, kafka-consumer-groups --describe shows - for those partitions, and lag exporters either drop them or report them as missing. Data ingestion is not affected.

Suggested fix: keep the dedup, but re-commit unchanged offsets once the last successful commit is older than a refresh interval set well below typical broker retention (for example one hour, or a dyncfg). Because #39461 switched to send_if_modified, process_frontier is never called for a fully idle source. So the refresh needs a timer in resume_uppers_process_loop, for example a select! over the stream and a tokio::time::interval, rather than only a looser condition here. The same timer would also retry a failed commit on an idle source. Today the doc comment's "retried on the next resume upper" may never happen for such a source.

martykulma and others added 2 commits October 2, 2026 10:20
Brokers expire a standalone group's offset per partition
offsets.retention.minutes after its last commit, so offsets the dedupe
skips must still be recommitted periodically. The interval is the
kafka_offset_commit_refresh_interval dyncfg, 10 minutes by default and
10 seconds in CI, with zero disabling the refresh.

Co-Authored-By: Claude <noreply@anthropic.com>
@martykulma
martykulma force-pushed the maz-kafka-offset-commit-dedupe branch from 7e97474 to e1893ed Compare October 2, 2026 14:21
@martykulma
martykulma requested a review from a team as a code owner October 2, 2026 14:21
@martykulma martykulma changed the title SS-455 - source: Kafka only commits offsets for a partition if they've changed SS-455 - source: commit Kafka offsets on change and on a timer Oct 2, 2026
@martykulma

Copy link
Copy Markdown
Contributor Author
  1. MEDIUM -- Idle partitions are never re-committed, so the broker expires their consumer-group offsets

Confirmed and addressed

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants