Repository navigation
Conversation
| -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" |
There was a problem hiding this comment.
❯ 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" # <---| /// 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, |
There was a problem hiding this comment.
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.
| let Some((batch_index, batch_size)) = message_id.batch() else { | ||
| self.entries.remove(&entry); | ||
| return Some(entry); | ||
| }; |
There was a problem hiding this comment.
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.
011651b to
b3fa0c9
Compare
|
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 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:
Worth noting that across connection resets, |
|
@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). |
|
@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. |
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
b3fa0c9 to
3ee9b4b
Compare
|
@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. |
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 ofevery message in the entry, processed or not.
This adds support for
BatchAcknowledgmentfollowing the design ofthe Java client.
ConsumerBuilder::with_batch_acknowledgmentselectshow a message inside a batched entry is acknowledged:
Entry, the default, acks the whole entry on the first ack of any ofits messages, which is the behavior described above and what every
earlier release did.
Trackedkeeps a bitset of unacked indexes per delivered entry andacks the entry once every message in it has been acked. This is the
Java client's behavior with
batchIndexAckEnabled=false, the defaultbefore Pulsar 4.1.
Indexedadditionally sends the remaining bitset as ack_set with eachack, so a broker with
acknowledgmentAtBatchIndexLevelEnabled=truerecords indexes itself and redelivers only the unacked ones. This is
the Java client's behavior with
batchIndexAckEnabled=true.Entryis kept as the default to avoid potentially-breaking changes forapplications. This is because
TrackedandIndexedhold an entryuntil 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 thatconsumer stops).
As a side-effect of this change we end up with better support for
cumulative acknowledgements. Under
TrackedandIndexed:entry (Tracked), or that plus the acked prefix of the entry (Indexed);
ack_setthe broker attaches to a redelivery is honored: themessages it reports as acked are skipped and their permits returned;
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
exclusiveorfailoversubscription makes the broker resend everything in flight.Closes #345
ci: apply the standalone broker configuration from env
bin/pulsar standalonedoes not read the PULSAR_PREFIX_* variables, sonone of the settings the workflow passes were actually applied.