Skip to content

[Issue 1534] Preserve batch size in transactional acknowledgments - #1535

Open
Ebispongebob wants to merge 1 commit into
apache:masterfrom
Ebispongebob:codex/fix-1534-transaction-batch-ack
Open

Ebispongebob wants to merge 1 commit into
apache:masterfrom
Ebispongebob:codex/fix-1534-transaction-batch-ack

Conversation

@Ebispongebob

Copy link
Copy Markdown

Fixes #1534

Motivation

A Shared consumer can stop receiving batched messages even though its transactions commit successfully. The Go client drops the producer batch size when converting a completed batch acknowledgment to an entry-level ID, then omits MessageIdData.BatchSize from the transaction ACK. Pulsar 4.0.3 therefore decrements the consumer's unacked count by one instead of the number of messages in the batch. The discrepancy accumulates until the broker blocks dispatch.

The issue includes a self-contained reproduction using the unmodified v0.21.0 client. With a broker unacked limit of 20 and three messages per producer batch, the 11th batch cannot be received after the previous ten transactions committed successfully.

Modifications

  • Preserve the producer batch size in completed batch acknowledgments and mark the resulting ID as whole-entry (batchIdx=-1), so ordinary ACK grouping does not reinterpret it as a partial batch ACK.
  • Populate the existing protobuf BatchSize field for transaction ACKs when the size is positive. Non-batched messages retain their existing wire behavior.
  • Use an independent outstanding-index set for each transactional batch-index ACK. Reusing the accumulated set would acknowledge earlier indexes again and cause transaction conflicts; the ACK must identify only the current message.
  • Add TestTxn_BatchAckUnackedMessages, using explicit producer flushes and broker stats to verify delivery continues and unacked/backlog return to zero. Cover batched and non-batched messages, single-message batches, ordinary ACKs, batch-index ACKs, abort/redelivery, and separate transactions within one producer batch.

Verifying this change

  • Make sure that the change passes the CI checks.

Local validation against official master 61d7a95e66cddf329adc6d925a8bd60ccd26eda5 plus this patch:

  • PASS: all eight TestTxn_BatchAckUnackedMessages subtests against Pulsar 4.0.3 standalone with transactionCoordinatorEnabled=true, acknowledgmentAtBatchIndexLevelEnabled=true, and maxUnackedMessagesPerConsumer=20.
  • PASS: existing consumer ACK, ACK grouping, duplicate detection, transaction-state error, and ACK tracker tests.
  • PASS: make lint (golangci-lint v2.12.2, 0 issues), go vet ./pulsar, and git diff --check.
  • The full TLS-dependent integration suite has not been run locally; CI remains unchecked until it completes.

Run the new regression with a transaction-enabled broker at localhost:6650 and its admin endpoint at localhost:8080:

go test ./pulsar -run '^TestTxn_BatchAckUnackedMessages$' -count=1 -timeout=4m

Enable broker-side batch-index acknowledgment for the batch-index subtests. The test does not change broker configuration. A low unacked limit reproduces the dispatch stall quickly; with the default limit, the final broker counter assertions still detect the accounting leak.

Does this pull request potentially affect one of the following parts:

  • Dependencies: No.
  • The public API: No.
  • The schema: No; the protobuf field already exists.
  • The default values of configurations: No.
  • The wire protocol: Yes. Transaction ACKs for batched messages populate the existing BatchSize field, and transaction batch-index ACK sets identify only the current index. No new fields or command types are introduced.

Documentation

  • Does this pull request introduce a new feature? No.
  • Feature documentation: Not applicable. This corrects existing transactional acknowledgment behavior and introduces no new public configuration.

Fixes apache#1534

Signed-off-by: 魏睿韬 <weiruitao@sofunny.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.

[Bug] Transactional ACKs omit batch size, causing Shared consumers to stop receiving

1 participant