[
https://issues.apache.org/jira/browse/QPID-8750?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Marco Geri updated QPID-8750:
-----------------------------
Summary: [Broker-J] AsyncAutoCommitTransaction can enqueue messages out of
order when a queue's messageDurability forces storage of non-persistent
messages (was: AsyncAutoCommitTransaction can enqueue messages out of order
when a queue's messageDurability forces storage of non-persistent messages)
> [Broker-J] AsyncAutoCommitTransaction can enqueue messages out of order when
> a queue's messageDurability forces storage of non-persistent messages
> --------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: QPID-8750
> URL: https://issues.apache.org/jira/browse/QPID-8750
> Project: Qpid
> Issue Type: Bug
> Components: Broker-J
> Affects Versions: qpid-java-broker-10.0.1
> Reporter: Marco Geri
> Priority: Major
>
> When a burst of non-persistent (transient) messages is published to a queue
> whose _messageDurability_ attribute is {_}ALWAYS{_}, the messages are
> occasionally enqueued in an order different from the order in which they were
> received on the connection.
> _AsyncAutoCommitTransaction_ decides whether to write to the message store
> using the queue's durability policy combined with the message's delivery mode:
> {code:java}
> if (queue.getMessageDurability().persist(message.isPersistent()))
> {
> txn = _messageStore.newTransaction();
> ...
> future = txn.commitTranAsync(null);
> } {code}
> but it then tells _addEnqueueFuture()_ whether the operation was persisted
> using the message's delivery mode alone:
> {code:java}
> addEnqueueFuture(future, new Action() { ... }, message.isPersistent()); {code}
> _addEnqueueFuture()_ relies on that flag to decide whether it is safe to
> invoke _postCommit()_ (which performs the actual enqueue onto the queue)
> synchronously:
> {code:java}
> private void addEnqueueFuture(final CompletableFuture<Void> future, final
> Action action, boolean persistent)
> {
> if (action != null)
> {
> // For persistent messages, do not synchronously invoke postCommit
> even if the future is completed.
> // Otherwise, postCommit (which actually does the enqueuing) might be
> called on successive messages out of order.
> if (future.isDone() && !persistent &&
> !_strictOrderWithMixedDeliveryMode)
> {
> action.postCommit();
> }
> else
> {
> _futureRecorder.recordFuture(future, action);
> }
> }
> } {code}
> The two values diverge whenever the queue's durability policy overrides the
> message's delivery mode. With _MessageDurability.ALWAYS_ and a non-persistent
> message, _MessageDurability.persist(false)_ returns true, so a store
> transaction is created and _commitTranAsync()_ returns a genuinely
> asynchronous future (for the BDB store, one completed later by the coalescing
> committer thread). However, persistent is passed as false, so
> addEnqueueFuture() takes the synchronous fast path as soon as the future
> happens to be already done.
> The result is a race: message N+1, whose commit future happens to be complete
> at the moment it is checked, is enqueued immediately, while message N, whose
> future is still pending, is deferred to the FutureRecorder (e.g.
> {_}AMQChannel._unfinishedCommandsQueue{_}) and enqueued only when that future
> completes. Consumers therefore see N+1 before N.
> Note that this is not the documented mixed-delivery-mode trade-off that
> _qpid.strict_order_with_mixed_delivery_mode_ addresses: here reordering
> occurs within a uniform stream of non-persistent messages, purely because of
> the queue's durability configuration.
> The same defect exists in both _enqueue()_ overloads of
> {_}AsyncAutoCommitTransaction{_}. The _Collection<? extends BaseQueue>_
> overload is the one used by the normal publishing path _(RoutingResult.send()
> → ServerTransaction.enqueue(Collection, ...)),_ so it is the more commonly
> hit of the two.
> The inverse mismatch is also wrong, though harmless: with
> _messageDurability=NEVER_ and a persistent message, no store write happens,
> but _persistent=true_ causes an already-completed future to be handed to the
> FutureRecorder, needlessly deferring _postCommit()_ until the next sync().
> Steps to reproduce:
> # Create a durable queue with messageDurability=ALWAYS on a virtualhost
> backed by the BDB message store (default coalescing-sync commit behaviour).
> # Publish a burst of non-persistent messages with a monotonically increasing
> sequence number in an application property, without waiting
> Occasionally the sequence is out of order. The frequency depends on store
> commit timing, so it is easier to reproduce under load or on slower storage.
> Suggested fix:
> Pass the value actually used to decide whether a store write was issued,
> rather than {_}message.isPersistent(){_}.
> For {_}*enqueue(TransactionLogResource, EnqueueableMessage,
> EnqueueAction)*{_}:
> {code:java}
> CompletableFuture<Void> future;
> final MessageEnqueueRecord enqueueRecord;
> final boolean persist =
> queue.getMessageDurability().persist(message.isPersistent());
> if (persist)
> {
> LOGGER.debug("Enqueue of message number {} to transaction log. Queue :
> {}",
> message.getMessageNumber(), queue.getName());
> txn = _messageStore.newTransaction();
> enqueueRecord = txn.enqueueMessage(queue, message);
> future = txn.commitTranAsync(null);
> txn = null;
> }
> else
> {
> future = CompletableFuture.completedFuture(null);
> enqueueRecord = null;
> }
> final EnqueueAction underlying = postTransactionAction;
> addEnqueueFuture(future, new Action() { /* unchanged */ }, persist); {code}
> For {*}_enqueue(Collection<? extends BaseQueue>, EnqueueableMessage,
> EnqueueAction)_{*}, the equivalent condition is whether a store transaction
> was opened for any of the target queues:
> {code:java}
> CompletableFuture<Void> future;
> final boolean persist = txn != null;
> if (persist)
> {
> future = txn.commitTranAsync(null);
> txn = null;
> }
> else
> {
> future = CompletableFuture.completedFuture(null);
> }
> final EnqueueAction underlying = postTransactionAction;
> addEnqueueFuture(future, new Action() { /* unchanged */ }, persist); {code}
>
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]