Skip to content

fixes #830: Close PC when InvalidPidMappingException - #839

Merged
Roman Kolesnev (rkolesnev) merged 2 commits into
confluentinc:masterfrom
mikhelef:close-pc-when-invalid-pid-mapping
Nov 4, 2024
Merged

Roman Kolesnev (rkolesnev) merged 2 commits into
confluentinc:masterfrom
mikhelef:close-pc-when-invalid-pid-mapping

Conversation

@mikhelef

@mikhelef Massinissa IKHELEF (mikhelef) commented Oct 31, 2024 •

Copy link
Copy Markdown
Contributor

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

  • Documentation (if applicable)
  • Changelog

@confluent-cla-assistant

confluent-cla-assistant Bot commented Oct 31, 2024 •

Copy link
Copy Markdown

🎉 All Contributor License Agreements have been signed. Ready to merge.
✅ mikhelef
Please push an empty commit if you would like to re-run the checks to verify CLA status for all contributors.

@mikhelef Massinissa IKHELEF (mikhelef) changed the title Close PC when InvalidPidMappingException fixes #830: Close PC when InvalidPidMappingException Oct 31, 2024
@mikhelef
Massinissa IKHELEF (mikhelef) force-pushed the close-pc-when-invalid-pid-mapping branch 8 times, most recently from 869a009 to 1c4a91a Compare November 1, 2024 09:20
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
@mikhelef
Massinissa IKHELEF (mikhelef) force-pushed the close-pc-when-invalid-pid-mapping branch from 015ffbf to 0c413b8 Compare November 1, 2024 11:00
@rkolesnev

Copy link
Copy Markdown
Contributor

/sem-approve

@rkolesnev

Copy link
Copy Markdown
Contributor

Looks good, thanks.
Couple of nitpicks:
Can you please update changelog - add section for 0.5.3.3 and move this fix under that section? Version 0.5.3.2 is already released - so no new changes should go into that section.

When changelog is updated - please run either ascii-doc maven target through IDE or mvn ascii-doc:template-build to update the Readme content with changelog update.

@mikhelef

Copy link
Copy Markdown
Contributor Author

Hi Roman Kolesnev (@rkolesnev)

I have added 0.5.3.3 section in changelog and updated readme.

@rkolesnev
Roman Kolesnev (rkolesnev) merged commit 1dbc015 into confluentinc:master 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

2 participants