Repository navigation
Keep instantiated OffsetMapCodecManager so that metrics will not be recreated every commit #859 - #892
Conversation
…ecreated every commit confluentinc#859 We can see the OffsetMapCodecManager being instantiated every second or so, which creates 3 new Meters (Timer, Gauge, Counter) every time. Because they are references from the registeredMeters array in PCMetrics, this creates a significant memory leak over time.
|
Nacho Muñoz Gómez (@nachomdo) or Roman Kolesnev (@rkolesnev) Are you maybe part of the reviewer group and can have a look at this PR? |
|
Stephen Hegarty (@stephenheg) John Byrne (@johnbyrnejb) Would any of you be available to review this fix? I was hoping to get the memory leak fixed before black friday :-) |
|
Martyn Ye (@sangreal) Are you part of the reviewer group and can help out here with a review? |
|
sorry, I am not a official Confluent reviewer |
Roman Kolesnev (rkolesnev)
left a comment
There was a problem hiding this comment.
LGMT.
Stefan Pries (@priesus) - Changelog / readme needs to be updated please as well - you can see how its updated in other PRs - its the CHANGELOG.adoc file and then once its updated - regenerate readme to take in the changes - mvn asciidoc-template:build .
John Byrne (@johnbyrnejb) - could you please do formal approval / merge to master once docs are updated?
Hey Stefan Pries (@priesus) - i dont work for Confluent anymore and dont use Parallel Consumer in my current day to day - but i try to keep an eye on it when i have a bit of time... |
|
Roman Kolesnev (@rkolesnev) Thanks so much for the feedback! I've updated the changelog & readme. |
|
John Byrne (@johnbyrnejb) I believe my changes broke the master branch build. I'm unable to see the build logs and I'm also unable to run the tests myself because our company policy doesn't let us use the maven repositories required by this project. |
|
Just to clarify: To verify my changes I had to copy the /src/main/java into one of our projects and I was able to fix the memory leak with the changes from this PR. So I assume the build issue is related to a test. Let me see if I can troubleshoot this somehow. |
|
John Byrne (@johnbyrnejb) So far I am still unable to get the tests to run locally, because I am unable to fetch the |
it wouldn't have bene broken, but the truth def is on githubs maven distribution system. I'm going to publish it to maven central soon to avoid this problem. |
…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>
…already carries confluentinc#892) Investigated the merged confluentinc#892 diff against the fork tree. confluentinc#892 hoists OffsetMapCodecManager to a single field on PartitionState (created once in the ctor) and is ALREADY in the fork -- on master, hence in PR #57s base. No conflict: our confluentinc#859 fix is layered on top (it uses #892s om field) and fixes the distinct leak source in PCMetrics.java (List -> LinkedHashSet + synchronized removal that prunes the tracking set). The authors "broke the master build" scare was the io.stubbs.truth dependency-availability issue, not a code regression -- the fork ships confluentinc#892 on its green master. Upstream confluentinc#859 stays open only because upstream merged confluentinc#892 but not our tracking-collection fix (status mixed); downstream it is fixed. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…f truth)
Add an optional per-entry backlink field to upstream-map.yaml. When set,
upstream-backlink.sh renders the upstream comment from it (with the same
{{FORK_REPO}}/{{FORK_REF}}/{{SUMMARY}}/{{ID}} placeholders) instead of the
generic template - so a tailored public explanation lives in the source of
truth, not a separate body file. Entries without the field still fall back to
the templates.
Use it on bug-859-pcmetrics-leak to explain the two-cause leak and how our
List->LinkedHashSet tracking fix complements the already-merged upstream confluentinc#892.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Render backlink comments from an optional per-entry backlink field in
upstream-map.yaml (source of truth) instead of a separate body, with the same
{{FORK_REPO}}/{{FORK_REF}}/{{SUMMARY}}/{{ID}} placeholders; entries without it
fall back to the generic templates. bug-859 uses it to explain the two-cause
leak vs the already-merged upstream confluentinc#892.
Make fork status honest about landed-ness. The old "fixed" conflated "fix
written" with "shipped": every "fixed" entry is actually an OPEN, unmerged fork
PR (or a branch with no PR). Replace with a lifecycle vocabulary
(none|in-progress|ready|pr-open|merged|released|superseded|wontfix) and correct
the entries: confluentinc#859/confluentinc#893/confluentinc#905 -> pr-open (in open PR #57), confluentinc#857 -> ready
(branch-only). Add an optional per-entry todo: list for outstanding actions
(merge the open PR, post the backlink) surfaced by "upstream-map.py todo", so
"still to do" is explicit rather than implied by an open PR + null forwarded.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…tus guard Audited PR #57 against issue confluentinc#859. The code matches the issue (registeredMeters List -> LinkedHashSet + prune, plus caching OffsetMapCodecManager in PartitionStateManager), but the wording framed the leak as rebalance-driven when the issue is commit-driven. Tighten bug-859 summary/notes/backlink: the List accumulated a duplicate Meter.Id on every registration (every commit); fork PR #57 fixes it via List->Set + PartitionStateManager caching (also closes #233); upstream confluentinc#892 covers the per-commit churn; confluentinc#893/confluentinc#905 are unrelated cherry-picks. Also fix upstream-backlink.sh: the fix-backlink status guard still checked the removed "fixed" value, so it refused every entry after the lifecycle-status change. Now allows ready|pr-open|merged|released and refuses none|in-progress. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
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>
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>
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>
…#72) Prep for the 0.6.0.0 release, from reconciling CHANGELOG.adoc against the 37 commits since upstream 0.5.3.3: - Fold the PCMetrics fix (upstream confluentinc#892 / confluentinc#859, "keep OffsetMapCodecManager instantiated so metrics aren't recreated every commit") into the 0.6.0.0 Fixes. It was filed under a phantom "0.5.3.4" heading, but 0.5.3.4 was never tagged - we go 0.5.3.3 -> 0.6.0.0 - so the fix actually ships in 0.6.0.0. Replace that heading with an explanatory comment. - Add a 0.6.0.0 Improvements line for the project rename to "Kafka Parallel Consumer" (branding only; Java package names unchanged, imports unaffected). - Regenerate README.adoc (it include::s CHANGELOG.adoc via the template). - release.yml: the GitHub-release step now uses our curated CHANGELOG section as the release notes instead of --generate-notes. It extracts the "== <version>" section, converts AsciiDoc to Markdown (=== -> ###, url[text] -> [text](url)), strips // comments, and passes --notes-file, falling back to --generate-notes if the section is missing/empty. Claude-Session: https://claude.ai/code/session_01T6Biwo4QXfD9CFrZnCq1zw Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Bring #57 current with master (release-notes #72, deps refresh #73/#74, unit-suite parallelisation #68, refactoring backlog #67, self-hosted CI, etc.). Only CHANGELOG.adoc conflicted: the 0.6.0.0 Fixes now lists master confluentinc#892 (per-commit OffsetMapCodecManager fix) alongside this PR confluentinc#859 (List->Set + assignment-path caching) and confluentinc#893 - complementary, kept all three. Regenerated README.adoc from the merged CHANGELOG via the asciidoc-template plugin. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…afe if encoding parallelises Surfaced while reviewing #57. Since confluentinc#892/#57 the codec manager is cached/shared (per-partition PartitionState.om for encode; one PartitionStateManager instance for decode). Correct today because encoding is single-threaded on the control thread, but encodingCounters (a plain HashMap) and the static errorPolicy would race if the #200 thread-model refactor ever parallelises encoding. Recorded under the OffsetMapCodecManager #233 entry as a refactoring-backlog item, not a bug. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

We can see the OffsetMapCodecManager being instantiated every second or so, which creates 3 new Meters (Timer, Gauge, Counter) every time.
Because they are referenced from the registeredMeters list in PCMetrics, the duplicate Meters and Tags won't be cleaned up by garbage collection. This creates a significant memory leak over time (issue #859).
This is my first open source contribution, please let me know if this change doesn't make sense or what contribution guidelines I have missed!
Checklist