Skip to content

[Feature] Add option to select batch construction strategy - #915

Open
Devingryu (devingryu) wants to merge 3 commits into
confluentinc:masterfrom
devingryu:features/batch-strategy
Open

Devingryu (devingryu) wants to merge 3 commits into
confluentinc:masterfrom
devingryu:features/batch-strategy

Conversation

@devingryu

Copy link
Copy Markdown

Adds the batchStrategy option to make batch construction behavior explicit. (See also: #266)

  • Adds batchStrategy option to ParallelConsumerOptions:
    • SEQUENTIAL: keeps the existing strict ordered behavior; in KEY/PARTITION modes only one record is taken per shard in a scheduling cycle, so shard-level concurrency remains effectively one-at-a-time. (Default)
    • BATCH_MULTIPLEX: allows taking multiple records from the same shard in one scheduling cycle, while still allowing a batch to contain records from multiple shards.
    • BATCH_BY_SHARD: allows taking multiple records from the same shard and additionally ensures that each batch contains records from only one shard.

Overview of the implementation

  • The work-selection path in ProcessingShard was changed so that, in ordered modes, the shard scan no longer always stops after the first taken record.
  • Previously, the loop effectively behaved as strictly sequential because it would stop scanning the shard immediately after a take.
  • With BATCH_MULTIPLEX and BATCH_BY_SHARD, the loop now continues gathering records from the same shard until batchSize is reached; SEQUENTIAL still takes the old stop-after-one path.
  • On the batch-construction side, BATCH_BY_SHARD uses shard-aware partitioning, so each batch is limited by batchSize and also split when the shard changes.

Concerns

  1. Semantics in UNORDERED
    UNORDERED + BATCH_MULTIPLEX currently behaves the same as size-based batching, while UNORDERED + BATCH_BY_SHARD means one topic-partition per batch. This is now documented and tested, but it would be good to confirm that this is the intended public API behavior.

  2. Ordered work selection vs. final batch construction
    The change splits behavior into two stages: how many records may be taken from a shard, and how the taken records are grouped into final batches. Please review whether this separation matches the intended model for the option.

Checklist

  • Documentation (if applicable)
  • Changelog

P.S. Since this repository does not appear to be actively maintained right now, I would particularly appreciate feedback from someone with a strong understanding of the project's internals.

* BATCH_MULTIPLEX allows batching to poll multiple messages from same shard.

* BATCH_BY_SHARD ensures batch to contain only messages from single shard.
@devingryu
Devingryu (devingryu) requested a review from a team as a code owner March 10, 2026 12:04
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Add a "Parallel-safe work while PR #57 is in flight" section to
docs/inflight.md recording, for each in-flight track, whether it
collides with PR #57's metrics/state files (857, 909, 51 -> sequence
after) or is parallel-safe (912, release, logging cleanup, security
bumps, contributor fixes, #40, confluentinc#915, DLQ), ranked by readiness.

Also refresh the confluentinc#859 entry to the consolidated PR #57 (bundles the
confluentinc#893/confluentinc#905 cherry-picks, supersedes the closed #42->#43->#45 stack) and
expand the confluentinc#912 entry (ready, pushed, no PR, vertx-isolated). Bump the
last-updated date.

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
Add a "Parallel-safe work while PR #57 is in flight" section to
docs/inflight.md recording, for each in-flight track, whether it
collides with PR #57's metrics/state files (857, 909, 51 -> sequence
after) or is parallel-safe (912, release, logging cleanup, security
bumps, contributor fixes, #40, confluentinc#915, DLQ), ranked by readiness.

Also refresh the confluentinc#859 entry to the consolidated PR #57 (bundles the
confluentinc#893/confluentinc#905 cherry-picks, supersedes the closed #42->#43->#45 stack) and
expand the confluentinc#912 entry (ready, pushed, no PR, vertx-isolated). Bump the
last-updated date.

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 Aug 17, 2026
Pointer note for the ideation doc on this branch: the open item is the
composition-API shape decision gating confluentinc#915, plus the
verified quick wins (over-request arithmetic fix) and the two confluentinc#915
blockers (PollContext ordering, missing batch telemetry/validation)
that should be raised before that PR is reviewed.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014ap9HhK65pD3qzdCVfBb25
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…lf-tuning from the backlog

The 2026-08-14..17 ideation wave produced decisions and directions the roadmap did not carry:
docs/data/roadmap.yaml still described the pre-proxy world, so the largest open workstream in the
repo (the language-proxy sidecar, #242, PR #293) appeared nowhere a reader could
watch it finish.

Six new entries at next-0x, which is the horizon the 1.0 exit criterion
intended-functionality-present inspects - putting them there is what records 'before 1.0':

- language-proxy-sidecar: the sidecar, the frozen v1 protocol, and clients answering to one
  conformance suite (#242, PR #293).
- polyglot-proof-demo: every binding visibly reading the same records, machine-checked.
- performance-comparison-matrix: scenarios x arms x languages measured in the conformance harness,
  published as data, under a fairness charter.
- distributed-throttling: divide a downstream budget across the group using the group itself
  (#228); throttling defers, liveness outranks the limit.
- self-tuning-concurrency: promoted from the backlog's 'dynamic concurrency control' entry
  (#227) - STRATEGY.md now carries this as the self-tuning bet, so the roadmap follows.
- batch-composition-settlement: choose the composition API once before confluentinc#915 freezes a
  shape into the options; answers #145 / confluentinc#266 in the same decision.

Three new backlog entries record what was considered and deliberately not scheduled: the proxy
HTTP dialect and REST Proxy compat gateway (feasibility settled, demand-gated), the Dapr component
(probe-gated, then demand-gated), and the actor/IPC revival (direction undecided; the probes and
the async-produce promotion #230 come first either way).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 19, 2026
… gate, and the self-tuning promotion (#314)

The roadmap now carries the 2026-08-14..17 ideation wave, says how real every entry is, and
enforces that it stays true.

Six new entries at next-0x - the horizon the 1.0 exit criterion inspects: the language-proxy
sidecar (#242), the polyglot proof demo, the performance comparison matrix, distributed
throttling (#228), self-tuning concurrency (#227, promoted from the backlog to
follow STRATEGY.md's self-tuning bet), and the batch-composition settlement (#145,
confluentinc#266) that must precede confluentinc#915. Three backlog entries record what was
considered and deliberately not scheduled: the proxy HTTP/REST-compat gateway, the Dapr component,
and the actor/IPC revival. The running-instance-visibility entry is renamed web-gui - the name
everything else uses.

Every entry now carries stage and stage_detail - an eight-value ladder (idea, ideated,
requirements-drafted, planned, limited-poc, poc, in-progress, implemented) ordered by the most
advanced artifact that exists, with definitions deliberately hard to flatter: implemented requires
proven use, and a brand-new system nobody has run in anger is a poc however large. An optional
stage_delivery records the readiness judgment (draft vs pending-merge) for entries carried by an
open PR. The schema makes stage/stage_detail entry_required, so no entry can claim a horizon
without saying how real the work behind it is - built for the v6 release announcement to draw
from directly.

The ladder's ownership rule - the PR advancing a track moves its entry's stage in the same change
- gets its own gate: roadmap-stage-gate.js compares the claiming entry's stage block between base
and head (editing another entry does not count), accepts only the fork's qualified astubbs#N
carrier form, and offers a reasoned roadmap-stage: N/A opt-out. Seventeen unit tests run before
the gate in CI, one pinned against the real roadmap file. A merge-checklist line covers entries
the gate cannot reach.

Also corrects the module-records note that pointed at git history the #273 squash made
unreachable - the drafted Connect feature record now rides its module-landing PR (#269) -
and merges master so the tree actually contains the STRATEGY.md self-tuning section the roadmap
cites (a review catch: the branch was cut six minutes before #308 merged).

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.

1 participant