Skip to content

feat(consumer): track acknowledgments of batched entries - #428

Open
mdeltito wants to merge 2 commits into
streamnative:masterfrom
mdeltito:batch-index-ack
Open

mdeltito wants to merge 2 commits into
streamnative:masterfrom
mdeltito:batch-index-ack

Conversation

@mdeltito

Copy link
Copy Markdown

feat(consumer): track acknowledgments of batched entries

Acknowledging one message of a batched entry sends a plain ack of the
entry's (ledger, entry) id, which the broker takes as an ack of
every message in the entry, processed or not.

This adds support for BatchAcknowledgment following the design of
the Java client. ConsumerBuilder::with_batch_acknowledgment selects
how a message inside a batched entry is acknowledged:

  • Entry, the default, acks the whole entry on the first ack of any of
    its messages, which is the behavior described above and what every
    earlier release did.
  • Tracked keeps a bitset of unacked indexes per delivered entry and
    acks the entry once every message in it has been acked. This is the
    Java client's behavior with batchIndexAckEnabled=false, the default
    before Pulsar 4.1.
  • Indexed additionally sends the remaining bitset as ack_set with each
    ack, so a broker with acknowledgmentAtBatchIndexLevelEnabled=true
    records indexes itself and redelivers only the unacked ones. This is
    the Java client's behavior with batchIndexAckEnabled=true.

Entry is kept as the default to avoid potentially-breaking changes for
applications. This is because Tracked and Indexed hold an entry
until it is fully acked, so an application that never (n)acks some
messages will hold back acks for the entries, which can end up hitting
the broker's maxUnackedMessagesPerConsumer (and thus dispatch to that
consumer stops).

As a side-effect of this change we end up with better support for
cumulative acknowledgements. Under Tracked and Indexed:

  • a cumulative ack in the middle of an entry acks up to the preceding
    entry (Tracked), or that plus the acked prefix of the entry (Indexed);
  • the ack_set the broker attaches to a redelivery is honored: the
    messages it reports as acked are skipped and their permits returned;
  • redelivery timeouts and negative acks are keyed by entry, as the
    broker redelivers whole entries.

In every mode, a cumulative ack now also stops the redelivery timer of
every entry it covers, which matches the Java client. Previously, only
the acked id was dropped, so the timers of already covered entries
still fired and requested redelivery, which on an exclusive or
failover subscription makes the broker resend everything in flight.

Closes #345


ci: apply the standalone broker configuration from env

bin/pulsar standalone does not read the PULSAR_PREFIX_* variables, so
none of the settings the workflow passes were actually applied.

-e PULSAR_PREFIX_acknowledgmentAtBatchIndexLevelEnabled=true \
apachepulsar/pulsar:${{ matrix.pulsar-version }} \
bin/pulsar standalone
bash -c "bin/apply-config-from-env.py conf/standalone.conf && bin/pulsar standalone"

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

❯ docker run --name probe -d -e PULSAR_PREFIX_advertisedAddress=9.9.9.9 \ # <---
          apachepulsar/pulsar:3.0.8 bin/pulsar standalone
  sleep 20
  docker exec probe curl -s http://127.0.0.1:8080/admin/v2/brokers/configuration/runtime \
          | tr ',' '\n' | grep '"advertisedAddress"'
  docker rm -f probe

"advertisedAddress":"localhost" # <---

❯ docker run --name probe -d -e PULSAR_PREFIX_advertisedAddress=9.9.9.9 \ # <---
          apachepulsar/pulsar:3.0.8 bash -c "bin/apply-config-from-env.py conf/standalone.conf && bin/pulsar standalone"
  sleep 20
  docker exec probe curl -s http://127.0.0.1:8080/admin/v2/brokers/configuration/runtime \
          | tr ',' '\n' | grep '"advertisedAddress"'
  docker rm -f probe

"advertisedAddress":"9.9.9.9" # <---

Comment on lines +15 to +18
/// Ack the whole entry on the first ack of any of its messages. The subscription moves
/// past the entry, unacked siblings included, and never delivers them again.
#[default]
Entry,

@mdeltito mdeltito Aug 24, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I kept this mode and also kept it the default, mostly to avoid a breaking change. There's an argument that, because this mode (and the existing behavior on master) is unsound and drops data, it should be removed entirely or at minimum not the default. Worth noting there is no equivalent to Entry for the Java client. It behaves like Tracked with batchIndexAckEnabled=false, and Indexed otherwise.

I think a reasonable path would be to cut over to Tracked as the default in a major release and potentially remove Entry.

Comment on lines +140 to +143
let Some((batch_index, batch_size)) = message_id.batch() else {
self.entries.remove(&entry);
return Some(entry);
};

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note that acks need valid batch_index and batch_size fields to be attributed to the batch. consumer.ack(&msg) and consumer.cumulative_ack(&msg) use the id the message was delivered with, so those work fine. But this means you cannot just rebuild the message ID with only (ledger_id, entry_id, batch_index) when using ack_with_id or cumulative_ack_with_id - you need to store all of the fields returned by msg.message_id(). For Tracked and Indexed, failing to provide the batch_size will fallback to the same old behavior as Entry.

`bin/pulsar standalone` does not read the PULSAR_PREFIX_* variables, so
none of the settings the workflow passes were actually applied.
@mdeltito

Copy link
Copy Markdown
Author

In addition to the new unit/integration tests, I've run a soak/chaos test locally against Pulsar 2.11.4, 3.0.8 and 4.1.3 and verified the behavior with batch-index setting on and off.

This test harness runs a producer/consumer pair through a TCP tap that decodes wire protocol frames with the crate's Codec. That lets it count acks, ack_set words, bytes per direction, batch sizes, and redeliveries from real frames instead of client logs. This also allows us to inject connection resets.

Every run ends by draining the subscription and comparing two sets against the exact set predicted by the mode and broker setting for that workload.

The tests varied by:

  • ack strategy: in order, shuffled, a fixed fraction of messages never acked, nack then ack on redelivery, ack the first delivery only, ack by an id rebuilt from its fields, and the new cumulative ack handling
  • subscription type: exclusive, shared, failover, key_shared. also tested with partitioned topics, unbatched producers, and varying batch sizes
  • faults: connection resets, topic unloads, consumer restarts (including restarting before pending acks were flushed, and changing mode on restart), one-off and periodic seeks, broker restarts, dead letter policies

Worth noting that across connection resets, Entry mode lost 98 messages the application never processed (because their entry was acked), which is the behaviour this PR exists to address. Tracked and Indexed lost none. The current master (v6.8.0 at the time) was also run against this harness and matched the behavior of Entry mode exactly.

@mdeltito

mdeltito commented Aug 25, 2026 •

Copy link
Copy Markdown
Author

@BewareMyPower @jiangpengcheng I'd love to get your thoughts on this, and a general sense of whether you would be open to accepting a PR for this improvement (this changeset, or in a different form).

@mdeltito

mdeltito commented Sep 8, 2026

Copy link
Copy Markdown
Author

@BewareMyPower @jiangpengcheng just checking in again here - this is an important feature for us, since the current state of the consumer with messages batched by the producer has the potential for significant data loss. I'm happy to adjust this in any way necessary so we can get support for indexed acknowledgement merged.

Comment thread src/consumer/engine.rs Outdated
Acknowledging one message of a batched entry sends a plain ack of the
entry's `(ledger, entry)` id, which the broker takes as an ack of
_every_ message in the entry, processed or not.

This adds support for `BatchAcknowledgment` following the design of
the Java client. `ConsumerBuilder::with_batch_acknowledgment` selects
how a message inside a batched entry is acknowledged:

- `Entry`, the default, acks the whole entry on the first ack of any of
  its messages, which is the behavior described above and what every
  earlier release did.
- `Tracked` keeps a bitset of unacked indexes per delivered entry and
  acks the entry once every message in it has been acked. This is the
  Java client's behavior with `batchIndexAckEnabled=false`, the default
  before Pulsar 4.1.
- `Indexed` additionally sends the remaining bitset as ack_set with each
  ack, so a broker with `acknowledgmentAtBatchIndexLevelEnabled=true`
  records indexes itself and redelivers only the unacked ones. This is
  the Java client's behavior with `batchIndexAckEnabled=true`.

`Entry` is kept as the default to avoid potentially-breaking changes for
applications. This is because `Tracked` and `Indexed` hold an entry
until it is fully acked, so an application that _never_ (n)acks some
messages will hold back acks for the entries, which can end up hitting
the broker's `maxUnackedMessagesPerConsumer` (and thus dispatch to that
consumer stops).

As a side-effect of this change we end up with better support for
cumulative acknowledgements. Under `Tracked` and `Indexed`:

- a cumulative ack in the middle of an entry acks up to the preceding
  entry (Tracked), or that plus the acked prefix of the entry (Indexed);
- the `ack_set` the broker attaches to a redelivery is honored: the
  messages it reports as acked are skipped and their permits returned;
- redelivery timeouts and negative acks are keyed by entry, as the
  broker redelivers whole entries.

In every mode, a cumulative ack now also stops the redelivery timer of
every entry it covers, which matches the Java client. Previously, only
the acked id was dropped, so the timers of already covered entries
still fired and requested redelivery, which on an `exclusive` or
`failover` subscription makes the broker resend everything in flight.

Closes streamnative#345
@mdeltito

Copy link
Copy Markdown
Author

@BewareMyPower @jiangpengcheng @david-streamlio @freeznet checking in here again to see if there is any additional feedback, and to see if this contribution might be something you all are willing to accept. Apologies for the repeated pings - I'm trying to plan around whether we can expect this functionality to be prioritized and incorporated, or if we need to go in a different direction and build this into applications directly.

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.

Batch Index Acknowledgment Support

2 participants