Skip to content

fix(state): trim noisy WARN in RemovedPartitionState to topic-partition + size + epoch - #918

Open
Vishal Kumar Singh (singhvishalkr) wants to merge 1 commit into
confluentinc:masterfrom
singhvishalkr:fix-631-removed-partition-state-log-level
Open

Vishal Kumar Singh (singhvishalkr) wants to merge 1 commit into
confluentinc:masterfrom
singhvishalkr:fix-631-removed-partition-state-log-level

Conversation

@singhvishalkr

Copy link
Copy Markdown

Fixes #631

Description

When a large number of topic-partitions are reassigned away from a consumer (e.g. during a rebalance), RemovedPartitionState#maybeRegisterNewPollBatchAsWork emits one WARN per dropped batch that dumps the entire EpochAndRecordsMap.RecordsAndEpoch:

https://github.com/confluentinc/parallel-consumer/blob/master/parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/RemovedPartitionState.java#L67-L69

// no-op
log.warn("Dropping polled record batch for partition no longer assigned. WC: {}", recordsAndEpoch);

Because RecordsAndEpoch carries the full List<ConsumerRecord<K,V>>, the resulting line can be many KB long. Common log tooling (Datadog, Splunk, CloudWatch) truncates at a fixed byte budget, which silently chops the most actionable part (topic-partition and epoch) out of the log.

Fix

Split the log into two statements, matching the pattern already used elsewhere in the project where detailed state is logged at DEBUG:

  • WARN: concise topic-partition (N records, epoch E) — short, grep-friendly, survives log-tooling truncation.
  • DEBUG (guarded by log.isDebugEnabled()): the full RecordsAndEpoch for deep troubleshooting when operators opt in.

The no-op behaviour, counts, and public API are completely unchanged; only the log output format is restructured.

Example before/after

Before (single ~5 KB line truncated to ~1 KB by log tool):

WARN Dropping polled record batch for partition no longer assigned. WC: EpochAndRecordsMap.RecordsAndEpoch(topicPartition=topic-7, epochOfPartitionAtPoll=42, records=[ConsumerRecord(topic = topic, partition = 7, leaderEpoch = 5, offset = 1234, CreateTime = 1708459200000, serialized key size = 12, serialized value size = 248, headers = RecordHeaders(headers = [], isReadOnly = false), key = k-1234, value = v-...[TRUNCATED]

After:

WARN Dropping polled record batch for partition no longer assigned: topic-7 (128 records, epoch 42)
DEBUG Full dropped batch for topic-7: EpochAndRecordsMap.RecordsAndEpoch(...)   ← only if debug is on

Checklist

  • Documentation (if applicable) — not applicable; log-format change only, no public API or user-facing doc references this line.
  • Changelog — this is a log-output restructuring; no behavioural change. Happy to add a CHANGELOG entry in whatever format the project prefers if a maintainer would like one.

Test plan

  • Existing RemovedPartitionState / partition-rebalance tests still pass; no behavioural assertions depend on the exact WARN text.
  • Manually verified the new format uses the Lombok @Value-generated getTopicPartition(), getRecords(), and getEpochOfPartitionAtPoll() accessors on the inner RecordsAndEpoch class in EpochAndRecordsMap.

…on+size+epoch

When many partitions are reassigned away from a consumer, the WARN log in
`RemovedPartitionState#maybeRegisterNewPollBatchAsWork` dumps the full
`EpochAndRecordsMap.RecordsAndEpoch` (including the entire ConsumerRecord
list) once per dropped batch. The resulting line is long enough to be
truncated by common log aggregators (Datadog, Splunk, etc.), which hides
the actually useful signal (which topic-partition was dropped, how many
records, at what epoch).

Split the log into:

- WARN: concise `topic-partition (N records, epoch E)` -- always readable
  and grep-friendly, survives log-tooling truncation.
- DEBUG (guarded by `log.isDebugEnabled()`): the full RecordsAndEpoch
  for deep troubleshooting when operators opt in.

Behaviour, counts, and no-op semantics are unchanged; only the log output
is restructured. This matches the approach already used elsewhere in the
codebase where detailed state is logged at DEBUG.

Fixes confluentinc#631
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
…p findings in manifest

The 2026-07-28 activity sweep surfaced upstream items the manifest didn't have.

Reconcile confluentinc#859: upstream PR confluentinc#892 (priesus) MERGED 2025-10-27 also targets the
PCMetrics leak, but fixes a DIFFERENT cause (OffsetMapCodecManager re-instantiated
each commit, recreating meters) than our fork fix (duplicate meter re-registration
on assign/revoke). Issue confluentinc#859 is still open, so mark upstream status 'mixed' and
add a reconciliation block tracking the open questions: does the fork already
carry confluentinc#892 / conflict with it, and did confluentinc#892 actually break the master build (the
author feared so; astubbs attributed it to the io.stubbs.truth dep not being on
Maven Central, not a code regression) -- verify before relying on it. Adds a
documented optional `reconciliation` field to the schema.

Capture other new findings: confluentinc#917 kafka-clients 3.9.2 SECURITY (into the security
batch); confluentinc#918/confluentinc#919 log-noise trims as a new logging-ux entry (with confluentinc#640/confluentinc#629/confluentinc#631);
confluentinc#920 JDK 17 build-doc PR linked to the Java-baseline entry; confluentinc#902 KEY-ordering
issue as a new entry; note confluentinc#921/confluentinc#922 (fork-awareness already partly upstream) on
the maintenance-signal entry so we don't duplicate it when backlinking.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Sep 1, 2026
…lure log lines

Two log lines interpolated a whole record batch, so each grew with
`max.poll.records` and with the assigned partition count until log tooling
truncated away the part the line existed to report.

- `RemovedPartitionState.maybeRegisterNewPollBatchAsWork` logged the entire
  `RecordsAndEpoch` at WARN (#169 / confluentinc#631). It now logs the
  topic-partition, the record count, the offset range and the epoch at poll -
  what an operator needs to tie the drop to a rebalance.
- The user-function failure line logged the whole `PollContextInternal`, which
  renders every record with keys and values (#170 / confluentinc#640). It
  now logs a summary of the same shape.

Both keep the unabridged object, one level down at DEBUG, so nothing is lost -
it just has to be asked for.

The formatting lives in one place,
`bz.stub.parallelconsumer.internal.utils.RecordBatchSummary`, rather than two
ad-hoc format strings: it names at most 5 partitions before collapsing the rest
to a count, so the line has a fixed ceiling no matter how big the batch or the
assignment. `AbstractParallelEoSStreamProcessor` is a contended file, so its own
diff is two lines (the logic sits behind `PollContextInternal#summariseForLog`).

The DEBUG line sits INSIDE `ThrowableUtils.logWithoutEscaping`, not after it.
Rendering the full `PollContextInternal` calls `toString()` on user keys and
values - user code, on the failure path, before the rethrow that the caller
depends on. That is precisely the hazard `logWithoutEscaping` exists to contain,
and the unabridged dump is the largest amount of user code this method runs.

Tests assert the emitted line, not the format string - "this line stays short"
is only true until someone interpolates a collection into it again, and only a
test that reads the log notices. New shared test helper `LogCapture` captures one
class's events and restores its level afterwards; `AmbientProbeExtensionTest` had
the same code inline and now uses it. Two further inline copies live in
`SubmitWorkToPoolShutdownRaceTest` and are deliberately left alone - unrelated to
these two issues, and recorded as a follow-up in the inflight note.

Both new assertions were proven to discriminate: reverting each production line
to its pre-fix form takes exactly its own test red ("expected to contain:
... 500 records, offsets 100-599", and the user-function equivalent), and
restoring takes them green.

`UserFunctionFailureLoggingTest` is `@Isolated`: it raises the level of a logger
shared by every processor in the JVM, and under JUnit thread parallelism the
first run of it both read other tests' lines and flooded the timing-sensitive
shutdown tests with DEBUG, failing `closeAfterSingleMessageShouldBeEventBasedFast`
and `queuedMessagesNotProcessedOrCommittedIfSubmittedDuringShutdown`. With the
isolation the core unit suite is green: 545 tests, 0 failures, 8 skipped.

The dropped-batch test is added to the `RemovedPartitionStateTest` that already
exists on master (the incomplete-offsets no-op contract) rather than to a class
of its own, and both capture sites filter on a topic name unique to the test, so
a concurrent test's dropped batch cannot land in the count.

Written independently of upstream PRs confluentinc#918/confluentinc#919, which
are unmerged and predate this fork's internals.

Upstream-Issue: confluentinc#631
Upstream-Issue: confluentinc#640
Upstream-PR: confluentinc#918
Upstream-PR: confluentinc#919
Forwarded: no
Applied-Upstream: no

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CMeraeBL2ycXVHEaTHM86W
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Sep 1, 2026
…owners

The note carried `post-merge: exempt-file - ... it is deleted when that PR lands`.
docs/inflight/AGENTS.md forbids exactly that:

    Never leave a "delete this when #NN merges" marker on master. The merge is
    exactly when nobody is looking here, so the marker outlives the work and the
    next reader inherits a stale note that reads as live.

Merging would have ADDED a file whose own text says the merge removes it, so the
marker ships and then sits there. Same defect #202 had, fixed the same way:
the directory's four outcomes say migrate what outlives the work first, and `git rm`
only once nothing left is both true and unowned.

MIGRATED, because none of it dies with the branch:

- The defect class - eight live instances, two named as worth doing next, and three
  dismissed with reasons - is SPLIT into its own note,
  docs/inflight/bug-unbounded-log-lines.md, rather than compressed into
  docs/refactoring.md.
  That file is explicitly for refactors "too small to deserve their own note, a line
  or two each"; this carries anchors, a decision about where the fix belongs (on the
  type, not the call site) and the dismissals, which are the part a line would lose.
  RecordBatchSummary's javadoc now points at it, so the next person writing a batch
  log line meets it from the code rather than by grepping.
- Every anchor in it was re-grepped rather than copied, and two claims were corrected
  in the move. #168 (confluentinc#629) asked for exactly the topic/partition/
  offset that ConsumerOffsetCommitter's error line carries, so a fix there must keep
  the identifiers and cap only the metadata string - the note previously read as
  though the whole map could be summarised away. And the examples module is not the
  "single WorkContainer, one key no value" case the three library modules are: CoreApp
  interpolates a raw ConsumerRecord, whose Kafka toString does print the value. It
  stays dismissed (one record, bounded, user-editable example code), but for a
  different reason than it was recorded under.
- The two inline `ListAppender` copies still in SubmitWorkToPoolShutdownRaceTest, with
  the standing follow-up to treat LogCapture as the only way to do this, go to
  docs/refactoring.md beside the other copy-pasted-test-helper entries. That one IS a
  line or two, so it belongs there and not in a note.
- The branch-specific flake lead moves into UserFunctionFailureLoggingTest's own
  javadoc, next to the @isolated it explains: without that annotation,
  closeAfterSingleMessageShouldBeEventBasedFast and
  queuedMessagesNotProcessedOrCommittedIfSubmittedDuringShutdown failed. Naming them
  there makes the lead reachable by grepping either test's name, which is how a future
  debugger of those two would arrive - a deleted branch note is not.

CONFIRMED ALREADY OWNED, so it dies with the note rather than being restated:

- The fork<->upstream mapping - the upstream-pr-log-noise entry in
  src/docs/development/upstream-map.yaml carries it, including why confluentinc#918/confluentinc#919 were read
  and not taken, and that the entry stays open after this lands.
- LogCapture's two hazards - its class javadoc states both, and which fix each takes.
- The log.debug inside logWithoutEscaping - the comment beside it, written in 394f462.
- The #201 overlap - settled: origin/fix/155-load-factor-noise has merged this
  branch and LoadFactorCeilingReportingTest now uses the shared LogCapture, so there is
  nothing left to coordinate in either merge order.
- "AbstractParallelEoSStreamProcessor.java is a contended file" - transient by nature,
  and the note itself said the answer is `gh pr diff --name-only`, not a file.

Nothing in the tree cited the deleted note: `grep -rn pr-log-verbosity` is empty, and
check-all.sh (15 gates, including check-file-refs, check-inflight-tags,
check-branch-self-reference and check-issue-refs) plus `bin/todo-index.sh --check` are
clean after the removal. Core unit suite unchanged at 546 run / 0 failures / 0 errors
/ 8 skipped.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CMeraeBL2ycXVHEaTHM86W
@astubbs

Copy link
Copy Markdown
Contributor

Vishal Kumar Singh (@singhvishalkr) - the maintained fork shipped this shape in v0.6.0.0: the dropped-batch WARN names the partition, the batch size and the epochs, and the full batch moved to DEBUG (astubbs#203, astubbs#428). Thank you for the report and the patch; the fork is where the work continues.

The coordinates changed with the fork - on Maven Central it is now:

<dependency>
    <groupId>bz.stub.parallelconsumer</groupId>
    <artifactId>parallel-consumer-core</artifactId>
    <version>0.6.0.0</version>
</dependency>

(parallel-consumer-vertx, parallel-consumer-reactor and parallel-consumer-mutiny likewise.) The Java package is bz.stub.parallelconsumer too, so imports change with it - one command rewrites them (GNU sed; on macOS use sed -i ''):

find . -name '*.java' -exec sed -i 's/io\.confluent\.parallelconsumer/bz.stub.parallelconsumer/g' {} +

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.

Warning log to verbose in RemovedPartitionState

2 participants