Repository navigation
fixes #830: Close PC when InvalidPidMappingException - #839
Merged
Roman Kolesnev (rkolesnev) merged 2 commits intoNov 4, 2024
Merged
Roman Kolesnev (rkolesnev) merged 2 commits into
Roman Kolesnev (rkolesnev) merged 2 commits into
Conversation
|
🎉 All Contributor License Agreements have been signed. Ready to merge. |
Massinissa IKHELEF (mikhelef)
force-pushed
the
close-pc-when-invalid-pid-mapping
branch
8 times, most recently
from
November 1, 2024 09:20
869a009 to
1c4a91a
Compare
Add fix in CHANGELOG.adoc Add fix in CHANGELOG.adoc add failure reason when closing due to invalidPidMappingException update UT with failurecause make failure reason setter protected refactoring refactoring
Massinissa IKHELEF (mikhelef)
force-pushed
the
close-pc-when-invalid-pid-mapping
branch
from
November 1, 2024 11:00
015ffbf to
0c413b8
Compare
Contributor
|
/sem-approve |
Contributor
|
Looks good, thanks. When changelog is updated - please run either ascii-doc maven target through IDE or |
Contributor
Author
|
Hi Roman Kolesnev (@rkolesnev) I have added 0.5.3.3 section in changelog and updated readme. |
Roman Kolesnev (rkolesnev)
approved these changes
Nov 4, 2024
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 2, 2026
#411 mirrors confluentinc#830, the only field report behind this work. Created on demand rather than by a sweep: the 2026-08-05 bulk import covered all 78 OPEN upstream issues and a separate cohort covered the 28 closed by the 2023 administrative sweep, and confluentinc#830 is in neither - it was closed in 2024 by a genuinely merged fix. Closed-and-actually-fixed issues were out of scope by design, on the reasoning that they have nothing outstanding. This one does: the fork intends to reverse the fix, which makes the original report load-bearing evidence. docs/inflight/upstream-coverage-completeness.md already registers the general obligation. confluentinc#839 deliberately gets neither a mirror nor a manifest entry. Mirrors are for issues, and the manifest keys on FORK work (fork.branches, fork.prs, fork_issue) - confluentinc#839 is an upstream PR that merged upstream and is already carried here as 1dbc015, so there is no fork-side mapping to record. It is covered inside the #411 body instead.
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 2, 2026
…quirements that rested on them Six-persona document review of the plan. Two findings were verified against the tree and both overturn something the plan asserted; the rest close gaps the plan left. WHAT WAS WRONG - "Aborting leaves the offsets dirty and the records are redelivered" is false for an instance that SURVIVES. PartitionState.onSuccess removes the offset from the incomplete set and raises the succeeded watermark at success time; onOffsetCommitSuccess only records the committed offset and clears the dirty flag. Redelivery today is a consequence of the instance DYING and rebuilding its offset state from the broker. Recovery in place would commit offsets whose output the abort discarded - silent result loss, which is precisely the doubt confluentinc#839 was settled on. The plan told the reader that doubt was answerable; it is not, without an explicit un-completion step, which is the pairing Kafka Streams makes between resetProducer and closeRunningTasksDirty. Now R13. - The produce-path catch does not fire on the ack wait. FutureRecordMetadata.valueOrError throws new ExecutionException(exception), so futureSend.get() can only surface an ExecutionException, which falls past the typed InvalidPidMappingException catch into the generic handler - PCInternalRuntimeException("Error while waiting for produce results", e), verbatim the stack trace in confluentinc#830. closePCWhenInvalidPidMappingException mocks a synchronous throw from produceMessages, so it passes without covering the reported path. Detection must unwrap first (now R9), and confluentinc#839's shutdown appears not to fire for the case it was written for - recorded as a defect in its own right rather than assumed handled. WHAT WAS MISSING - Producer configuration would print at INFO on every startup: ParallelConsumerOptions carries a bare @tostring and is logged whole; today the field is a Producer whose toString is an identity hash, a config map is SASL and keystore secrets. Now R7. - The derived transactional.id had no ACL story - substituting it stops an operator's existing TransactionalId grant matching, and initTransactions fails fatally. Now R6 (documented stable group-derived prefix) and R19 (re-grant as an upgrade step). - Failure of the recovery attempt itself was unspecified, so the natural implementation hot-loops on the control thread through a coordinator outage. Now R14. - The single response omitted the group rejoin, so the group-generation member of the condition set recurs immediately and never converges. Now in R10. - The transactional.id had to be stable and reused, or the replacement never fences its predecessor and the old transaction blocks read_committed consumers until it times out. Now in R4. - Unbounded recovery had no non-progress signal, so a permanently fencing instance looks healthy. Now R22, inside the settled logs-and-metrics boundary. - The configuration path was scoped to transactional flows while the deprecation covered every flow, stranding non-transactional producers. R1 widened; recovery stays transactional-only. Requirements renumbered R1-R22; every Governs, Covers and Trigger reference updated. Four new acceptance examples. One open question raised rather than settled: whether the deprecated producer-instance option's removal is queued for the major being cut now.
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 2, 2026
…e the roadmap entry it belongs to TWO THINGS THE #225 REVIEW SURFACED THAT ITS OWN PLAN DELIBERATELY DOES NOT FIX bug-411-wrapped-send-failure-spins-forever.md. ParallelEoSStreamProcessor catches InvalidPidMappingException around the produce-and-ack block and closes PC. That catch was confluentinc#839, written to end the infinite retry loop #411 (confluentinc#830) reported. It cannot fire for that report: FutureRecordMetadata.valueOrError throws new ExecutionException(exception), so the ack wait can only surface an ExecutionException, which falls past the typed catch into the generic handler producing PCInternalRuntimeException("Error while waiting for produce results", e) - verbatim the stack trace in the report, from a build that already carried the fix. From there the record is marked failed and re-dispatched onto the same invalid producer, which is the reported spin. The typed catch fires only on a synchronous throw from produceMessages, which is exactly what closePCWhenInvalidPidMappingException mocks - so the test is green over a path the reporter never took. Not fixed by the #225 work: that plan requires unwrapping before matching, but only where PC builds its own producer, and every user today supplies a Producer instance. The plan excludes it in Scope Boundaries rather than absorbing it, and now cites this note as the owner. Two ways out are on the page; choosing between them is a product call, so the note records both rather than assuming one. Deliberately NOT claimed: that the spin reproduces on master today. The reasoning is from the current tree and the kafka-clients 3.9.2 sources; nobody has run it. The note says what a reproduction needs. ROADMAP survive-producer-fencing advances idea -> requirements-drafted. Its stage_detail said "no design yet", which stopped being true when the requirements plan landed, and product review flagged the entry as underselling the banked design. done_when is rewritten to what the contract actually commits to - a replaced producer, no offset from the aborted transaction committed, its work processed again - and to say the guarantee applies where PC builds the producer, since the deprecated path keeps today's behaviour. Adds the field report and the PR as related links.
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 2, 2026
… unwrap it from the shape it really arrives in
One recoverable condition, found wherever kafka-clients puts it.
RecoverableProducerCondition.find walks a failure's cause chain for
ProducerFencedException, InvalidProducerEpochException,
InvalidPidMappingException, OutOfOrderSequenceException (so also
UnknownProducerIdException) and CommitFailedException, through the
wrappers the client uses: ExecutionException from a send future,
KafkaException("...we are in an error state") from every later transactional
call, and PC's own PCInternalRuntimeException. confluentinc#839's catch matched
the outermost type, which is why it fired on the synchronous shape and never
on the one in the field report (#411, confluentinc#830).
Where PC built the producer (canRecover), a match is recorded on the
ProducerManager - first condition wins, any thread may record, nobody blocks -
and the operation unwinds with ProducerInvalidatedException: a worker's record
fails and returns through the mailbox, a commit unwinds through the existing
lock release. Recovery itself is the control thread's, and lands next. On the
deprecated producer-instance path every condition keeps its pre-recovery
outcome, pinned by tests on both the commit path and all four produce-path
shapes - including the wrapped-send-future spin that
docs/inflight/bug-411-wrapped-send-failure-spins-forever.md owns, now visible
in the suite rather than only in a note.
Two things inherited from #262 are fixed while here, both in the code
this touches. The instance-path close on InvalidPidMappingException swallowed
the exception, so the batch was marked succeeded for output never produced;
it now rethrows after the close, and a test with the close stubbed out shows
the batch failing and retrying rather than being produced once and committed
(mutation-checked: restoring the swallow fails it). And a throwing abort on
close skipped closing the producer; abort is now swallowed and the close
happens regardless.
ProducerManager gains the replacement-supplier constructor PCModule wires in;
the four-argument constructor stays for the tests that build one directly.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01VrpH51xNDodaajE4P2nhFg
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 3, 2026
…ose never leaks the producer (#429) Two exactly-once holes on the transactional produce path, found by the review of #262 and recorded there rather than fixed. Both closed here, one commit's worth of reasoning each. Hole 1: ParallelEoSStreamProcessor#processAndProduceResults caught InvalidPidMappingException, called closeOnException, and fell through to return a partial - usually empty - result list, so runUserFunctionInternal marked EVERY WorkContainer in the batch succeeded. Output never produced, offset advanced: the same shape as the partial-result-set defect #261 fixed. It now rethrows after the close, wrapped as a PCInternalRuntimeException exactly as the generic arm does, so the batch is FAILED and returned to the mailbox and its offsets can never become commit payload. The instance still closes with the same failure cause. Why nobody saw it: on master the wrong verdict never reached a commit by accident, because closeOnException blocks the worker until the instance is CLOSED and the succeeded containers land in a mailbox nobody drains. The verdict was wrong regardless, and #410 is redesigning exactly the blocking close that masked it, which is why this had to land first. Field history: confluentinc#830 reported the spin, confluentinc#839 wrote this arm; the reporter's wrapped ExecutionException shape takes the generic arm and is #411, deliberately untouched here. Hole 2: ProducerManager#close(Duration) aborted an open transaction inside try/finally that released the lock but did not contain the abort's throw, with closeProducer(timeout) after the block. A fenced or poisoned producer's abort always throws, so close() leaked the KafkaProducer - IO thread, sockets, buffers - while doClose still marked the instance CLOSED. The whole cleanup is now contained, the commit lock is released only if it was taken, and closeProducer runs from a finally. The defect-class sweep found two siblings, both fixed: the commit-lock timeout path released a lock it never took (IllegalStateException "Not held be me", same leak, different door), and innerDoClose's producer step was the one shutdown step without the try/catch(WARN) its three neighbours carry from confluentinc#818, so a throwing producer close overwrote failureReason with "Error from poll control thread" and ran doClose twice. Checked and dismissed: commitOffsets releases through AbstractOffsetCommitter's own finally; RetryQueue's lock pairs take the lock in the statement before the try; brokerPollSubsystem.drain() sets a field and calls the one consumer method deliberately left unguarded, and is recorded in the write-up rather than guarded without evidence. RED then GREEN, four tests, each observed failing before its fix: invalidPidMappingFailsTheBatchInsteadOfMarkingItSucceeded (three predicted assertions all observed: the record never re-dispatched, the offset committed, the record marked done; control arm re-run after the review round changed the test), closeStillClosesTheProducerWhenTheAbortThrows, closeStillClosesTheProducerWhenTheCommitLockCannotBeAcquired, and closeCompletesWhenTheProducerCloseThrows. The register gains no new claim - there is no documented sentence for it - so the hole-1 test carries @ProvesClaim(OFFSET_AND_RECORDS_ATOMIC) as a second arm of C4, recorded in C4's note. Rejected: unwrapping ExecutionException into the pid arm (changes the behaviour #411 tracks and collides with #410's recovery); a non-blocking signal-then-throw close (rewrites closeOnException's contract, which #410 owns); catching only around abortTransaction with the lock release unconditional (#410's current shape, which leaves the lock-timeout leak in place). The local review round found the hole-1 test leaking a non-daemon control thread into a shared Surefire fork - the spy that stubs closeOnException is a different object from the one the base @AfterEach closes - now closed in a finally. Retires docs/inflight/bug-eos-swallowed-produce-failures.md; the defect class both holes share is written up in docs/solutions/logic-errors/a-catch-that-closes-and-continues-reports-the-batch-succeeded-2026-09-03.md and the note's two citations are repointed. #410 must reconcile when it next merges master; the list is on that PR. The local branch fix/producer-close-commit-lock carried an older hole-2 fix and is superseded. Also records a RegistrationRaceStaleResidentIT setup-guard flake seen twice on this branch and once elsewhere, with a deliberate re-run of the same head passing at normal length, and attests the inherited LatencyUtils backlog line naming this PR. Upstream-Issue: confluentinc#830 Upstream-PR: confluentinc#839 Applied-Upstream: no Co-authored-by: Claude Fable 5.1 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #830
Problem
In transactional mode, when ParallelConsumer encounters an InvalidPidMapping exception, it retries records indefinitely leading to same exception.
Proposed Solution
Parallel Consumer should exit when encoutering InvalidPidMapping.
Checklist