Skip to content

fix(kafka): stop deferred offset stores from queuing behind heartbeat polls - #129

Merged
Divyanshu Tiwari (divyanshu-tiwari) merged 2 commits into
mainfrom
fix/kafka-store-offset-lock
Sep 23, 2026
Merged

Divyanshu Tiwari (divyanshu-tiwari) merged 2 commits into
mainfrom
fix/kafka-store-offset-lock

Conversation

@divyanshu-tiwari

Copy link
Copy Markdown
Contributor

Description

In group read mode, Finish keeps polling (heartbeatPoll, 1s Poll under consumerMu) while deferred acks settle. Each settled ack calls storeOffset, which also took consumerMu, and the ack tracker runs one ack at a time. With partitions paused, every poll blocks for its full second, so stores drained at only a few per second.

Runs that read tens of thousands of records therefore spent hours in Finish after downstream tasks (e.g. the file writer's _SUCCESS) had already completed. If an external timeout killed the process first, very little was committed. Because acks settle out of order, the stored position per partition barely advances. The next run re-read an even larger backlog.

StoreOffsets and Pause are thread-safe in librdkafka, so storeOffset and pausePartition no longer take consumerMu. The lock still serializes Poll/ReadMessage/Commit/Close; Close runs only after the tracker has drained, so no store can overlap it.

Validated with a local (uncommitted) test: 500 deferred acks settled under an active heartbeat loop went from ~19s to <1s. go test -race ./internal/pkg/pipeline/... passes.

Types of changes

  • Docs change / refactoring / dependency upgrade
  • Bug fix (non-breaking change which fixes an issue)
  • New feature (non-breaking change which adds functionality)
  • Breaking change (fix or feature that would cause existing functionality to change)

Checklist

  • My code follows the code style of this project.
  • My change requires a change to the documentation and I have updated the documentation accordingly.
  • I have added tests to cover my changes.

@divyanshu-tiwari
Divyanshu Tiwari (divyanshu-tiwari) merged commit b4a9c32 into main Sep 23, 2026
7 checks passed
@divyanshu-tiwari
Divyanshu Tiwari (divyanshu-tiwari) deleted the fix/kafka-store-offset-lock branch September 23, 2026 14:20
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.

3 participants