Conversation
8d00f5b to
6da992c
Compare
|
One coordination note for future contributions: I opened #26499 and indicated that I was willing to submit a PR, although I had not opened one yet. In this situation, it would be helpful to leave a short comment before starting work, so contributors can avoid duplicating effort. No problem this time—just please coordinate on the issue first in similar cases. One change is required: the new I did not find another blocker in the current diff. |
| return batchDeletedIndexes == null | ||
| || batchDeletedIndexes.size() <= getConfig().getMaxBatchDeletedIndexToPersist(); |
There was a problem hiding this comment.
Could we clarify the intended policy here? The method isCursorDataFullyPersistable() is used directly by multiple-consumer dispatchers when dispatcherPauseOnAckStatePersistentEnabled is enabled, so this change effectively makes managedLedgerMaxBatchDeletedIndexToPersist a dispatcher-pause threshold.
PIP-299 explicitly described pausing when managedLedgerMaxUnackedRangesToPersist is reached. While pausing on batch-index overflow also seems reasonable, this represents a real behavior change—not merely an improvement to the helper's accuracy.
If this is intended, could you update the PR description accordingly and add a dispatcher-level test that covers pause/resume behavior when the batch-index limit is exceeded?
|
|
||
| @Test | ||
| public void testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit() throws Exception { | ||
| ManagedLedgerConfig config = new ManagedLedgerConfig(); |
There was a problem hiding this comment.
After rebasing onto master, should we use rawEntryConfig() here?
This test writes raw byte arrays using ledger.addEntry(new byte[] {...}), and the existing tests in ManagedCursorBatchAckRecoveryTest have already migrated to rawEntryConfig() to prevent the cache/read path from parsing them as Pulsar message entries.
Keeping new ManagedLedgerConfig() here would make this new test inconsistent with the current setup and may cause issues after the rebase.
| long[][] ackSets = {{2L}, {Long.MIN_VALUE, 1L}, {5L, 0L, 3L}}; | ||
| for (int i = 1; i < positions.size(); i++) { | ||
| Position position = positions.get(i); | ||
| cursor.delete(AckSetStateUtil.createPositionWithAckSet( | ||
| position.getLedgerId(), position.getEntryId(), ackSets[i - 1])); | ||
| } |
There was a problem hiding this comment.
buildBatchEntryDeletionIndexInfoList() can persist at most maxBatchDeletedIndexToPersist records. It should therefore be verified that, with a limit of two, two batch-deleted-index records are reported as fully persistable, and that adding a third switches the result to false.
This also protects the check against a future accidental regression of >= to > (or vice versa).
…stable Also compare batchDeletedIndexes size to maxBatchDeletedIndexToPersist so dispatchers see non-persistable ACK state when batch records exceed the limit. Assisted-by: Cursor (Composer)
Align managed-ledger test with defaultConfig(), fix import order, and assert persistability at the batch-index cap boundary. Add a broker test that pauses and resumes dispatch when batch deleted index state exceeds the limit under PIP-299 dispatcher pause policy.
|
Thanks for the review — addressed in tip @void-ptr974 — Moved the
Tests run locally (Temurin 21):
|
6da992c to
c7457d0
Compare
|
|
||
| Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING).topic(tpName) | ||
| .subscriptionName(subscription).subscriptionType(SubscriptionType.Shared) | ||
| .receiverQueueSize(1).enableBatchIndexAcknowledgment(true).isAckReceiptEnabled(true) |
There was a problem hiding this comment.
I think there is still a test race condition around the consumer's initial Flow.
ConsumerImpl completes subscribeFuture before calling increaseAvailablePermits(cnx, getCurrentReceiverQueueSize()). After subscribe() returns, either the test's afterAckMessages() call or the connection thread sending the initial permit can win.
If the initial Flow wins, blockedDispatcherOnCursorDataCanNotFullyPersist remains false, and the dispatcher can already read one of the existing backlog messages. Since cancelPendingRead() cannot retract a message already delivered to the consumer queue, pausedReceive could become non-null even though the pause behavior is correct.
Could we make this deterministic by avoiding the initial permit? For example, using a zero-queue consumer with a pending receiveAsync() would allow us to set the blocked state first, then explicitly issue a single permit through the receive. The same future could be checked for timeout while paused and completed after the dispatcher resumes.
…ataFullyPersistable-batch-index Signed-off-by: arimu1 <19286898+arimu1@users.noreply.github.com>
Use receiverQueueSize(0) and a pending receiveAsync so afterAckMessages can set the blocked state before any permit is issued, addressing the subscribe/Flow race noted in review. Signed-off-by: arimu1 <19286898+arimu1@users.noreply.github.com>
Denovo1998
left a comment
There was a problem hiding this comment.
Thanks for the update. I checked the zero-queue approach against the consumer Flow path, and this addresses the race I was concerned about.
With receiverQueueSize(0), no initial Flow is sent during subscribe, and the pending receiveAsync() explicitly issues the permit only after the dispatcher has entered the paused state. Reusing the same future after the batch-index state becomes persistable also covers the resume path cleanly.
The previous review comments look addressed now. LGTM.
Fixes #26499
Motivation
ManagedCursorImpl.isCursorDataFullyPersistable()only compared the whole-entrydeletion range count to
maxUnackedRangesToPersist. When partial batch ACK stateexceeded
maxBatchDeletedIndexToPersist, the method still returnedtrueeventhough
buildBatchEntryDeletionIndexInfoList()truncates on persist. Thatmismatch can leave dispatchers running while some batch ACK records would not
survive a cursor-state write.
Modifications
isCursorDataFullyPersistable()to returnfalsewhenbatchDeletedIndexesexceedsmaxBatchDeletedIndexToPersist.ACK records with limit 2), including boundary checks at the cap.
dispatcherPauseOnAckStatePersistentEnabledis on, multi-consumer dispatchersalso treat
managedLedgerMaxBatchDeletedIndexToPersistas a pause thresholdvia
isCursorDataFullyPersistable()(in addition tomanagedLedgerMaxUnackedRangesToPersist).Verifying this change
This change added tests and can be verified as follows:
./gradlew :managed-ledger:checkstyleTest :managed-ledger:test --tests org.apache.bookkeeper.mledger.impl.ManagedCursorBatchAckRecoveryTest.testIsCursorDataFullyPersistableReflectsBatchDeletedIndexLimit./gradlew :pulsar-broker:test --tests org.apache.pulsar.client.api.SubscriptionPauseOnAckStatPersistTest.testPauseOnBatchDeletedIndexLimitExceededDoes this pull request potentially affect one of the following parts: