Skip to content

fix maxConcurrency documentation in javadoc - #650

Closed
Vinaya Prasad N (nvinayshetty) wants to merge 461 commits into
confluentinc:masterfrom
nvinayshetty:master
Closed

Vinaya Prasad N (nvinayshetty) wants to merge 461 commits into
confluentinc:masterfrom
nvinayshetty:master

Conversation

@nvinayshetty

Copy link
Copy Markdown

Description...

Checklist

  • Documentation (if applicable)
  • Changelog

Antony Stubbs (astubbs) and others added 30 commits July 23, 2021 17:57
This messes with the shutdown process of the Producer, i.e. committing final transactions, offsets etc.
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.
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
* 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
@nvinayshetty Vinaya Prasad N (nvinayshetty) changed the title fix maxConcurrency related in javadoc fix maxConcurrency documentation in javadoc Nov 3, 2023
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>
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.