Skip to content

fix(offset): don't block PartitionOffsetManager.Close without auto-commit - #3753

Open
hsdfat wants to merge 1 commit into
IBM:mainfrom
hsdfat:fix-pom-close-deadlock-2772
Open

hsdfat wants to merge 1 commit into
IBM:mainfrom
hsdfat:fix-pom-close-deadlock-2772

Conversation

@hsdfat

@hsdfat hsdfat commented Sep 19, 2026

Copy link
Copy Markdown
Contributor

With Consumer.Offsets.AutoCommit.Enable = false, PartitionOffsetManager.Close() never returns unless something else calls Commit() or Close() on the OffsetManager. The repro from the issue is ManagePartition followed by pom.Close().

pom.Close() calls AsyncClose() and then ranges over pom.errors. That channel is only closed by pom.release(), which runs from releaseSelectedPOMs. Without auto-commit, mainLoop isn't started, so only an explicit om.Commit() or om.Close() reaches it. The OffsetManager.Close doc says to close the POMs first, and doing that deadlocks. The existing manual-commit tests close the OffsetManager before the POM to work around it.

Now, when auto-commit is disabled, pom.Close() releases the POM from its parent itself. Two details:

  • The release runs in another goroutine while Close drains pom.errors. A concurrent Commit() can be blocked sending an error to a full pom.errors (with Consumer.Return.Errors) while holding pomsLock for reading. Taking the write lock before draining would deadlock, and draining first unblocks it. Close still waits for the release to finish before it returns.
  • The new releasePOM only releases the POM if it is still the one managing its partition. So calling Close again on a POM that was already released can't release a newer POM for the same partition, which releasing by topic/partition would do.

Behaviour change: in manual-commit mode, pom.Close() releases the POM without committing its marked offset, so callers need to Commit() first. This is what removePartitions does for revoked partitions and what OffsetManager.Close() does for the remaining POMs in this mode. Before, the only way Close returned was when another goroutine's Commit() or Close() released the POM. I added one sentence about this to the Close doc on the interface. Auto-commit mode is unchanged.

I considered and rejected:

  • Flushing before releasing. That would make Close send a commit the user didn't ask for, and neither of the other manual-mode release paths does that.
  • Releasing in AsyncClose. It is called with pomsLock read-held from asyncClosePOMs and removePartitions, so taking the write lock there would deadlock. It would also drop an offset that a later Commit() is meant to flush.

The consumer group session never calls pom.Close(); it uses removePartitions and om.Close(). release() is guarded by a sync.Once, and the POM is removed from the map under the write lock, so a later Commit/Close doesn't see it again. pom.lock isn't held when pomsLock is taken, which keeps the existing pomsLock → pom.lock order.

TestPartitionOffsetManagerCloseWithoutAutoCommit has three cases. Each fails with a 5s timeout on main:

  • It manages two partitions, marks an offset on one and closes it. Close returns, the closed partition is no longer managed while the other still is, and no OffsetCommit is sent.
  • With ChannelBufferSize = 1 and Return.Errors, one failing Commit() fills pom.errors and a second one blocks sending its error. The test waits until that commit holds pomsLock, then checks that Close returns both errors. Releasing before draining deadlocks here.
  • It closes a POM, manages the same partition again, AsyncCloses the new POM, and closes the old one a second time. The new POM is still managed.

go test -race . passes, apart from TestClientBackgroundMetadataUpdater and TestPartitionConsumerComputeBackoff, which fail the same way on main when the whole package runs (a leftover goroutine logging through the swapped Logger). golangci-lint run reports 0 issues.

Fixes #2772

…mmit

With Consumer.Offsets.AutoCommit.Enable set to false, Close on a
PartitionOffsetManager waited for its errors channel to be closed, but
only the parent's Commit or Close releases a closed POM, and without
auto-commit the parent has no loop calling Commit. Closing the POM
before the OffsetManager, as the OffsetManager.Close doc asks, blocked
until something else called Commit or Close on the OffsetManager.

Release the POM from its parent in Close when auto-commit is disabled.
Like removePartitions and OffsetManager.Close in this mode, it does not
commit the marked offset, which remains the caller's job via Commit.
The release runs in another goroutine while Close drains the errors,
because a concurrent Commit may be blocked sending an error while
holding pomsLock. It only releases the
POM if it still manages its partition, so closing it twice can't
release a newer POM for the same partition.

Fixes IBM#2772

Signed-off-by: phatlc <phatle.hsd@gmail.com>
@hsdfat
hsdfat force-pushed the fix-pom-close-deadlock-2772 branch from cab9c38 to 2736904 Compare September 21, 2026 06:41

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.

PartitionOffsetManager.Close deadlock without autocommit

1 participant