Repository navigation
fix maxConcurrency documentation in javadoc - #650
Closed
Vinaya Prasad N (nvinayshetty) wants to merge 461 commits into
Closed
Vinaya Prasad N (nvinayshetty) wants to merge 461 commits into
Vinaya Prasad N (nvinayshetty) wants to merge 461 commits into
Conversation
This messes with the shutdown process of the Producer, i.e. committing final transactions, offsets etc.
…ncy issues under high pressure
refactor: Extract common Reactor and Vert.x parts
Prevents the extension modules from incorrectly inheriting core methods that would be broken to use. Step 1: new parent: rename
Remove deprecated test class
Base class refactor - removes core api from extension modules
Don't know how this made it through CI. Missing copyright and updated readme.
… see #maxCurrency Vert.x concurrency control previously relied on Vert.x WebClient controlling concurrency setting per host. This breaks things when you use multiple hosts - no the max concurrency can go beyond the setting. This change migrates to the new ExternalEngine system - which controls concurrency properly. Turn off performance comparison unit test - too brittle for CI.
… test-jar dependency handling bug
CI will catch it.
Partitions now track their highest succeed offset, for use bye encoders. Instead of encoding only up to the highest succeeded offset, it would attempt to encode the entire buffered partition state. This along with a race condition in OffsetEncodingBackPressureTest#backPressureShouldPreventTooManyMessagesBeingQueuedForProcessing was causing the test on CI to fail intermittently (and sometimes locally). This change also makes the test behave consistently regardless of environment performance.
…ntinc#537) * Add Invalid Offset Metadata error policy option * Fix file headers * Fix tests * Add some javadocs to the new invalidOffsetMetadataPolicy option * Add suggested missing test
Feature: Metrics support through micrometer --------- Co-authored-by: Antony Stubbs <antony.stubbs@gmail.com> Co-authored-by: Nacho Munoz <nachomdo@gmail.com>
fix: example in readme, license headers
* PL-176 handle close dont drain mode gracefully
…remove add remove staled work containers function in ProcessingShard (confluentinc#623) * 1. exclude stale containers from counting 2. add remove staled work containers function in ProcessingShard 3. update changelog * ignore stale work containers while processing * remove stale containers when partition assigned and revoked * remove line * revert filter stale worker when invoking getCountOfWorkAwaitingSelection
* fix: Return cached pausedPartitionSet Signed-off-by: Taeik Lim <sibera21@gmail.com> * update license headers * Change to cache partition count only Signed-off-by: Taeik Lim <sibera21@gmail.com> --------- Signed-off-by: Taeik Lim <sibera21@gmail.com> Co-authored-by: Edward Vaisman <10497078+eddyv@users.noreply.github.com>
…#619) Signed-off-by: Taeik Lim <sibera21@gmail.com> Co-authored-by: Roman Kolesnev <88949424+rkolesnev@users.noreply.github.com>
confluentinc#627) * add synchronization to ensure proper intializaiton and closing of PCMetrics singleton. fixes confluentinc#617
* fix readme doc
* PL-468 - Refactor metrics to support multiple PC instances in same java process
…onfluentinc#626) Bumps [io.projectreactor:reactor-core](https://github.com/reactor/reactor-core) from 3.5.7 to 3.5.9. - [Release notes](https://github.com/reactor/reactor-core/releases) - [Commits](reactor/reactor-core@v3.5.7...v3.5.9) --- updated-dependencies: - dependency-name: io.projectreactor:reactor-core dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
Bumps [org.xerial.snappy:snappy-java](https://github.com/xerial/snappy-java) from 1.1.10.1 to 1.1.10.3. - [Release notes](https://github.com/xerial/snappy-java/releases) - [Commits](xerial/snappy-java@v1.1.10.1...v1.1.10.3) --- updated-dependencies: - dependency-name: org.xerial.snappy:snappy-java dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
confluentinc#612) Bumps `junit.platform.version` from 1.10.0-RC1 to 1.10.0. Updates `org.junit.platform:junit-platform-launcher` from 1.10.0-RC1 to 1.10.0 - [Release notes](https://github.com/junit-team/junit5/releases) - [Commits](https://github.com/junit-team/junit5/commits) Updates `org.junit.platform:junit-platform-commons` from 1.10.0-RC1 to 1.10.0 - [Release notes](https://github.com/junit-team/junit5/releases) - [Commits](https://github.com/junit-team/junit5/commits) --- updated-dependencies: - dependency-name: org.junit.platform:junit-platform-launcher dependency-type: direct:development update-type: version-update:semver-patch - dependency-name: org.junit.platform:junit-platform-commons dependency-type: direct:development update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
…entinc#603) Bumps [maven-site-plugin](https://github.com/apache/maven-site-plugin) from 4.0.0-M8 to 4.0.0-M9. - [Release notes](https://github.com/apache/maven-site-plugin/releases) - [Commits](apache/maven-site-plugin@maven-site-plugin-4.0.0-M8...maven-site-plugin-4.0.0-M9) --- updated-dependencies: - dependency-name: org.apache.maven.plugins:maven-site-plugin dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
fix: array key equality - enables byte array key comparison to work correctly in Key ordered mode
Charles Provencher (cprovencher)
force-pushed
the
master
branch
from
November 3, 2023 18:14
075b4f9 to
558d599
Compare
4 tasks done
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Aug 12, 2026
…he sweep tool a mode that could have caught it (#258) * docs(upstream-map): track the 28 issues closed by the 2023 admin sweep Upstream ran two bulk clearouts before going quiet, and neither was a triage: 2023-06-15 `eddyv` closed 35 unmerged PRs with "Closing - Stale." (34 of them ours), and 2023-07-07 `johnbyrnejb` closed 28 issues with "Closing Issue" - every one marked COMPLETED rather than "not planned". That state reason is why the cohort was invisible. GitHub renders `completed` as resolved, so the issues read as done at a glance, and `upstream-sweep.sh` searches `updated:>=last_swept`, so anything last touched in 2023 can never appear in a sweep. None of the 28 were in this manifest; the mirroring that produced the other 78 fork mirrors only ever sampled *open* upstream issues. All 28 were re-read in full - bodies plus all 58 comments - and each verified against fork source rather than trusted from its thread. That mattered: four issue bodies point at PRs that "fix" them (#372->#390, #319->#270, #203->#345, #191->#346) and every one of those PRs was itself swept unmerged. Result: 2 already fixed (#41 offset-scan removal, #319 shutdown CME), 6 partly addressed, 20 fully open. Grouped into work items rather than 28 near-duplicate entries, since several want designing together - per-topic handlers with separate consume/produce types, seek with subscription changes, retry expiry with stall detection. #57 folded into the existing log-noise entry. Swept PR head commits were checked and are all still reachable, so the drafts are recoverable; entries cite branch and SHA, not just the PR number. Upstream-Issue: confluentinc#154 Forwarded: no Applied-Upstream: no * docs(upstream-map): mirror the whole 2023 cohort, not the parts we rate The point of mirroring is a clean, known cutoff - "every issue upstream closed administratively is accounted for here" - not a curated pick of the good ones. A set filtered by our own value judgement is not a cutoff: the next person cannot tell "not mirrored because it was junk" from "not mirrored because it was missed", and that ambiguity is worth far more than the cost of carrying a few items nobody will ever action. Drops the decline recommendations from #53, #199, #246 and #57 in favour of neutral statements of the trade-off, and renames sweep-2023-decline-candidates to sweep-2023-small-items (status wontfix -> none). The facts that informed those recommendations are kept - they are useful to whoever picks the issue up - but they are now framed as input to that person's decision rather than as our verdict delivered in advance. Fork issues can always be closed later on their merits, which leaves a visible record; pruning before mirroring does not. Records the principle in the section header so a later pass does not re-prune. Upstream-Issue: confluentinc#154 Forwarded: no Applied-Upstream: no * docs(upstream-map): link the created fork mirrors back to their entries The 28 mirrors are live as #227-254 (label `upstream-admin-closed`), with confluentinc#41 -> #233 and confluentinc#319 -> #252 created and immediately closed as completed, carrying the code evidence that they are genuinely fixed. Records the fork numbers on each entry so the mapping is queryable rather than re-derived, which is the reason this file exists. Adds `fork.fork_issues` (the plural of the existing `fork_issue`) because several entries deliberately group upstream items that want designing together, and documents it in the schema block. Verified: all 28 mirrors are referenced, none missing, no strays. Also snapshots each upstream title verbatim into its mirror body header, dated 2026-08-07 - the fork titles may drift as these are worked, and the original wording is what makes an old thread findable. Upstream-Issue: confluentinc#154 Forwarded: no Applied-Upstream: no * fix(tooling): give upstream-sweep a mode that can see closures it never saw The default sweep is a window search - `updated:>=$since` - so it structurally cannot see an item whose last activity predates the window, and `last_swept` only moves forward, so the blind spot grows. That is why upstream's 2023 administrative closures went unnoticed for three years: 28 issues closed with "Closing Issue" and 35 unmerged PRs closed as "Closing - Stale.", all marked COMPLETED, none of it triage, none of it visible to us. Adds `--audit`: no time window, asks "which closed upstream items are neither tracked in the manifest nor mirrored in the fork?", and flags days where an implausible number of things closed at once. Run against the real repo it rediscovers both sweeps from scratch and reports zero unaccounted, which is also an end-to-end check that the new fork_issues linkage is correct. Bots are excluded from the PR analysis. Dependabot self-closes superseded bumps in batches that look exactly like a sweep - 2022-08-16, 2022-10-20, 2023-11-03 and 2024-01-25 are all dependabot - and unfiltered they buried the two real sweeps in noise. Filtering them surfaced two human PRs that had been sitting inside those batches (confluentinc#508, confluentinc#650), now recorded for review. The report says outright that a bulk day is not proof of a sweep, because a release triage looks identical from here, and that stateReason COMPLETED is not evidence of a fix - that assumption is what made this cohort invisible. Known remaining hole, recorded not fixed: a PR closed alone on a quiet day is still invisible, since detection keys on bulk. Upstream-Issue: confluentinc#154 Forwarded: no Applied-Upstream: no * fix(tooling): teach the audit about Discussions, a content type we never checked Discussions were invisible to everything: not in the manifest, not mirrored, not in any sweep mode. 74 of them upstream, and nothing we run would ever have mentioned one. They were NOT swept - no bulk closure, and the largest day (2024-04-02, six threads) is all answered=true housekeeping. So the failure mode is different from the 2023 cohort: not administrative closure, just questions nobody answered, in a place no issue search reaches. Handling differs accordingly - these want answering or converting, not mirroring wholesale. The audit now lists zero-reply discussions, excluding release-announcement threads by title: those legitimately have no replies, and counting them as neglect would put 9 false positives at the top of the report. Two finds justify the check on their own. Discussion 542 (zero replies) is a field report of a transactional consumer stuck in rebalance because it cannot acquire the produce lock, with the reporter suspecting the revoke-time flush - the same lock lifecycle as bug-producing-lock-double-release, and they raised the timeout without effect, which is evidence about the mechanism rather than the duration. Discussion 883 (zero replies) is the same complaint as fork mirror #187 but with a runnable reproducer, which #187 lacks. Upstream-Issue: confluentinc#154 Forwarded: no Applied-Upstream: no * docs(inflight): record the open obligation to account for every upstream item Two notes, both about what is NOT done. Discussions are not a queue to be converted. There is no pipeline turning threads into issues, and most should never become one. The rule is a judgement: if reading a discussion makes us think there is an issue, we raise one - a normal fork issue on its own merits, citing the discussion as where it came from. It exists because we believe the problem is real, not because a thread existed. That also keeps `upstream-mirror` meaning exactly "an upstream issue we carry", which is what makes the 2023 cutoff verifiable. The wider obligation is now written down: every upstream issue, PR and discussion must be accounted for - carried, declined, or genuinely resolved upstream - not sampled. The 2023 cohort was found only because it was a *bulk* event; `--audit` keys on bulk closures and zero-reply threads, and both are proxies that miss whole shapes of problem (a PR closed alone on a quiet day, a discussion with one dismissive reply, an issue marked COMPLETED with no linked PR). The audit narrows the field; only reading discharges the obligation. Also records what has been ruled out, so it is not re-investigated: wiki disabled, no advisories, all open-milestone issues already mirrored, the orphan branches accounted for (v0.6.x-dev is the lambda-actor work already captured via swept PRs; 0.5.3.x's regression fix IS on master as a908e16), and "upstream pushed today" being a pushed_at artefact rather than new activity. Left open: project boards need a read:project scope we do not have, and 169 forks are unexamined. * fix(upstream-map): use a group the schema actually allows Two new entries used `group: batching-ordering`, which is not in the GROUPS allow-list in scripts/upstream-map.py, so `upstream-map.py validate` failed on them. Caught in review on #258. My mistake was validating the wrong way: I checked the file with an ad-hoc yaml.safe_load plus a duplicate-id check, which passes happily on a group the schema rejects, instead of running the project's own validator that AGENTS.md tells contributors to run before committing manifest changes. Reassigned both to `features` rather than widening GROUPS. The rest of this cohort's feature entries already use `features`, so this keeps the cohort internally consistent, and the ordering/batching distinction is not lost - it lives where it belongs, on the issues themselves via the `area/batching-ordering` label (#244, #236). Adding a vocabulary term to a shared schema as a side effect of a mirroring PR is the larger, less reversible change. Also drops the redundant 345 from sweep-2023-broker-disconnect-commit's `related` list, where it already appears in `prs`. Verified: `upstream-map.py validate` reports OK on 27 entries, and `refs` and `table` both still render. * docs(upstream): graduate the durable parts of the sweep record out of inflight The three inflight notes stay - they hold undischarged work - but they were also carrying permanent content, and inflight is transient by its own charter. Moved to durable homes, each citing the manifest rather than restating it: - docs/upstream.md gains the upstream-admin-closed cohort note (with the "do not trust 2023-era closure states" rule), the discussions non-mirroring policy decided 2026-08-07, documentation of --audit and its known blind spots, and the ruled-out upstream surfaces (wiki, advisories, milestones, orphan branches, pushed_at). - docs/solutions/ gains the transferable lesson: a closure state is a rendering choice not a triage, "fixed by #N" must be checked against the merge bit, windowed watchers need a no-window audit, and bots must be filtered before hunting bulk events. - upstream-pr-analysis.adoc is corrected: it claimed issues were never bulk-closed, which the 2023-07-07 sweep disproves, and the PR sweep is now confirmed rather than suspected. The inflight notes now hold only open work, pointing at the new homes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * docs(upstream): name the analysis doc's role (the plan) beside the manifest (the state) docs/upstream.md's opening enumerated everything it owns except the editorial analysis, and neither source-of-truth reference was a clickable link. Now the two roles are named explicitly: the .adoc is the plan, the manifest is the state tracker. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(upstream): make the audit's discussion query actually paginate, and correct the sweep cohort Address PR review feedback (#258). - upstream-sweep.sh: `gh --paginate` only substitutes a variable named `endCursor`; calling it `$cursor` meant the cursor was never fed back and gh re-requested page 1 forever (reproduced: 469 identical pages at first:2). It only looked healthy because 74 discussions fit in one page of 100. The failure was also swallowed by `2>/dev/null || true`, which would have printed a clean "(none)" for a broken query - the one answer this mode must never invent. Fail loudly instead. - upstream-map.yaml: the sweep cohort listed 36 PRs against a documented count of 35. Upstream closed exactly 35 unmerged PRs on 2023-06-15; confluentinc#66 (closed 2022-10-19) was carried in by mistake. Listing a ref here marks it accounted for, so it was hiding an unrelated PR from every future audit. Recorded PR 258 and moved the entry to `pr-open` per the lifecycle rule. - next-upstream-coverage-completeness.md: dropped the cached counts - inflight notes must not record what a command can answer, and the "~100 partially mirrored" line already contradicted docs/upstream.md, which says all 78 open upstream issues are mirrored with backlinks. Kept only what the audit cannot say: what nobody has read yet. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs(agents): let the manifest cache frozen upstream-issue facts, and say why The rule said "the manifest tracks upstream PRs only", so mirroring the 2023 sweep cohort into it read as a violation. The rule is what is out of date. An archived upstream's closed issue numbers and closure events are frozen. Caching them locally is a read-path optimisation, not a second tracker: grepping one file is instant, while the same answer from the mirrors costs dozens of API round-trips and burns rate limit shared across agents. The usual objection to duplication - the copy silently diverges - needs a source that can still move, and this one cannot. So the boundary is redrawn where it actually bites: the mirror owns an issue's *live state* and remains the only place you update it; the manifest may cache what is frozen (number, cohort, owning mirror). Cache what is frozen; never mirror what is moving. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs(upstream): put the manifest-caching rule in the doc that owns it The previous commit wrote the rule into AGENTS.md, which is the one place it does not belong. AGENTS.md's own preamble calls itself a router - "if it only matters once you are already in a topic... it goes in that topic's doc" - and docs/upstream.md line 7 already claims ownership explicitly: "AGENTS.md carries only the pointer and the one-line rule that manifest upkeep is the agent's job." The rule was in fact already stated twice before this branch (AGENTS.md:442 and docs/upstream.md:23), so growing the AGENTS.md copy deepened a fragmentation that was already there. Now stated once, where an agent editing the manifest will actually be looking, with AGENTS.md left holding the pointer and the one binding rule whose failure is silent (keep the fork side in sync). Net effect on the router: one line shorter than master, despite covering more. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs(agents): make a pointer say who owns the rule, not just what is nearby Diagnosed from a live miss. A reviewer cited AGENTS.md's "the manifest tracks upstream PRs only", so that section was read and treated as authoritative - while docs/upstream.md had been stating the same rule, and claiming ownership of it, the whole time. The duplicate was then grown in the router rather than fixed at source. Nothing in the reading path revealed the mistake. AGENTS.md stated a complete-looking rule and followed it with a pointer listing four adjacent mechanics ("the manifest schema, the mirrors, the commit trailers and the upstream sweep"), which reads as *further detail*, not as *this is a summary and that doc wins*. The contract was written down - docs/upstream.md:7 - but only on the side nobody enters from, since AGENTS.md is what loads every session. The anti-duplication rule already existed here ("Never state a fact twice"). What was missing is the half that makes it discoverable: a cross-reference must name the owner and say the owner wins, so a reader can tell a stub from the whole rule before editing the copy. Applied to the rule itself and to the upstream pointer that failed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <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.
Description...
Checklist