From dc97de1ef6234cc56844f4b9029c496381c5bb40 Mon Sep 17 00:00:00 2001 From: litiliu <38579068+litiliu@users.noreply.github.com> Date: Thu, 27 Aug 2026 16:19:21 +0800 Subject: [PATCH] [client] Retry stale out-of-order responses instead of resetting the idempotent writer An idempotent writer could be fatally reset by a stale OutOfOrderSequence response, killing all in-flight writes across every bucket. Failure scenario addressed: with max-inflight > 1 and acks=all, PutKv/ProduceLog responses are returned in completion order (a bucket's error result rides along with its request, whose response is gated by the slowest replicating bucket in that request). A batch (seqN) can transiently fail (e.g. NotLeaderOrFollower during a leader/partition transition) so its successor (seqN+1) hits an empty or behind writer state and is rejected with OutOfOrderSequence. seqN is then retried and acknowledged, advancing lastAckedBatchSequence, before seqN+1's already-generated out-of-order response is delivered. At that point the batch looks like the next expected sequence, so canRetry treated it as an unrecoverable regression and reset the writer id, which then dropped other in-flight batches with UnknownWriterId. Fix: snapshot each batch's lastAckedBatchSequence at send time (WriteBatch.lastAckedSequenceAtSend, set in Sender before dispatch, refreshed on every send attempt). In IdempotenceManager.canRetry, retry an OutOfOrderSequence when the acked sequence has advanced since the batch was sent (current lastAcked > lastAckedSequenceAtSend): that advancement means a predecessor was acknowledged while the batch was in flight, so the response is a superseded (stale) one that a reordering surfaced late, and a retry will observe the advanced state and succeed. A genuine regression shows no such progress and still resets. Adds SenderTest cases for the stale (retry, no reset) and genuine (reset) out-of-order paths, and one asserting the send-time snapshot is refreshed on retry so a repeated out-of-order response with no new ack still resets. --- .../client/write/IdempotenceManager.java | 12 +- .../org/apache/fluss/client/write/Sender.java | 17 +++ .../apache/fluss/client/write/WriteBatch.java | 12 ++ .../apache/fluss/client/write/SenderTest.java | 139 ++++++++++++++++++ 4 files changed, 179 insertions(+), 1 deletion(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/IdempotenceManager.java b/fluss-client/src/main/java/org/apache/fluss/client/write/IdempotenceManager.java index 1bc8f2bc12a..69a290d73f7 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/IdempotenceManager.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/IdempotenceManager.java @@ -269,7 +269,10 @@ synchronized boolean canRetry(WriteBatch batch, TableBucket tableBucket, Errors if (error == Errors.OUT_OF_ORDER_SEQUENCE_EXCEPTION && (batch.sequenceHasBeenReset() - || !isNextSequence(tableBucket, batch.batchSequence()))) { + || !isNextSequence(tableBucket, batch.batchSequence()) + || lastAckedBatchSequence(tableBucket) + .orElse(IdempotenceBucketEntry.NO_LAST_ACKED_BATCH_SEQUENCE) + > batch.lastAckedSequenceAtSend())) { // We should retry the OutOfOrderSequenceException if the batch is not the next batch, // i.e. its batch sequence isn't the lastAckedBatchSequence + 1. However, if the first // in flight batch fails fatally, we will adjust the batch sequences of the other @@ -279,6 +282,13 @@ synchronized boolean canRetry(WriteBatch batch, TableBucket tableBucket, Errors // OutOfOrderSequence, we want to retry it. To account for the latter case, we check // whether the sequence has been reset since the last drain. If it has, we will retry it // anyway. + // + // We should also retry if the acked sequence has advanced since this batch was sent + // (lastAckedBatchSequence > lastAckedSequenceAtSend). That means a predecessor was + // acknowledged while this batch was in flight, so this OutOfOrderSequence reflects a + // superseded (stale) server state that a response reordering surfaced late; a retry + // will observe the advanced state and succeed. A genuine loss shows no such progress + // and still resets. return true; } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/Sender.java b/fluss-client/src/main/java/org/apache/fluss/client/write/Sender.java index b16780035f7..f4f0c6eedf5 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/Sender.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/Sender.java @@ -375,6 +375,23 @@ private void sendWriteRequest(int destination, short acks, List return; } + // Snapshot each idempotent batch's send-time lastAckedBatchSequence, so a later + // out-of-order response can be recognized as stale (and retried) if the acked sequence + // has advanced while the batch was in flight. + if (idempotenceManager.idempotenceEnabled()) { + for (ReadyWriteBatch batch : batches) { + if (batch.writeBatch().hasBatchSequence()) { + batch.writeBatch() + .setLastAckedSequenceAtSend( + idempotenceManager + .lastAckedBatchSequence(batch.tableBucket()) + .orElse( + IdempotenceBucketEntry + .NO_LAST_ACKED_BATCH_SEQUENCE)); + } + } + } + // group record batch by table id. Map> writeBatchByTable = new HashMap<>(); batches.forEach( diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteBatch.java b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteBatch.java index 754f9caa706..9cbfaa6af2e 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteBatch.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteBatch.java @@ -57,6 +57,10 @@ public abstract class WriteBatch { protected boolean reopened; protected int recordCount; private long drainedMs; + // Snapshot of the bucket's lastAckedBatchSequence taken when this attempt was sent, used to + // tell a stale/reordered OutOfOrderSequence response (state advanced while in flight) apart + // from a genuine one (no progress since send). + private int lastAckedSequenceAtSend = IdempotenceBucketEntry.NO_LAST_ACKED_BATCH_SEQUENCE; public WriteBatch( long tableId, @@ -173,6 +177,14 @@ public boolean sequenceHasBeenReset() { return reopened; } + public void setLastAckedSequenceAtSend(int lastAckedSequenceAtSend) { + this.lastAckedSequenceAtSend = lastAckedSequenceAtSend; + } + + public int lastAckedSequenceAtSend() { + return lastAckedSequenceAtSend; + } + public int bucketId() { return bucketId; } diff --git a/fluss-client/src/test/java/org/apache/fluss/client/write/SenderTest.java b/fluss-client/src/test/java/org/apache/fluss/client/write/SenderTest.java index 69d9c669c9e..3e99500c317 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/write/SenderTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/write/SenderTest.java @@ -26,6 +26,7 @@ import org.apache.fluss.config.MemorySize; import org.apache.fluss.exception.AuthorizationException; import org.apache.fluss.exception.NetworkException; +import org.apache.fluss.exception.OutOfOrderSequenceException; import org.apache.fluss.exception.TableNotExistException; import org.apache.fluss.exception.TimeoutException; import org.apache.fluss.metadata.PhysicalTablePath; @@ -752,6 +753,144 @@ void testCorrectHandlingOfOutOfOrderResponsesWhenResponseLostButSubsequentBatche assertThat(sender1.numOfInFlightBatches(tb1)).isEqualTo(0); } + /** + * A late (reordered) out-of-order response must be retried, not treated as fatal, when the + * acked sequence has advanced while the batch was in flight. + * + *

Scenario: seq0 and seq1 are both in flight (each sent when {@code lastAcked=-1}). Response + * reordering delivers seq0's success first, advancing {@code lastAcked} to 0. seq1's + * out-of-order response then arrives late; it reflects a superseded server state. Since seq1 + * was sent when {@code lastAcked=-1}, which is now behind the current {@code lastAcked=0}, the + * writer must retry seq1 instead of resetting. + */ + @Test + void testStaleOutOfOrderResponseIsRetriedInsteadOfResettingWriter() throws Exception { + IdempotenceManager idempotenceManager = createIdempotenceManager(true); + Sender sender1 = setupWithIdempotenceState(idempotenceManager); + sender1.runOnce(); + assertThat(idempotenceManager.isWriterIdValid()).isTrue(); + long writerId = idempotenceManager.writerId(); + + // Send seq0 and seq1; both are sent while lastAcked is not present (-1). + CompletableFuture future0 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(1, "a"), (tb, leo, e) -> future0.complete(e)); + sender1.runOnce(); + CompletableFuture future1 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(2, "b"), (tb, leo, e) -> future1.complete(e)); + sender1.runOnce(); + assertThat(idempotenceManager.nextSequence(tb1)).isEqualTo(2); + assertThat(idempotenceManager.lastAckedBatchSequence(tb1)).isNotPresent(); + + // Reordered responses: seq0 succeeds first, advancing lastAcked to 0. + finishIdempotentProduceLogRequest(0, tb1, 0, createProduceLogResponse(tb1, 0L, 1L)); + sender1.runOnce(); + assertThat(future0.get()).isNull(); + assertThat(idempotenceManager.lastAckedBatchSequence(tb1)).isEqualTo(Optional.of(0)); + + // seq1's out-of-order response arrives late. It is stale (seq1 was sent at lastAcked=-1, + // now behind lastAcked=0), so the writer must retry it, not reset. + finishIdempotentProduceLogRequest( + 1, tb1, 0, createProduceLogResponse(tb1, Errors.OUT_OF_ORDER_SEQUENCE_EXCEPTION)); + sender1.runOnce(); + assertThat(idempotenceManager.isWriterIdValid()).isTrue(); + assertThat(idempotenceManager.writerId()).isEqualTo(writerId); + assertThat(future1.isDone()).isFalse(); + + // The retried seq1 succeeds once the server observes the advanced state. + sender1.runOnce(); + finishIdempotentProduceLogRequest(1, tb1, 0, createProduceLogResponse(tb1, 1L, 2L)); + sender1.runOnce(); + assertThat(future1.get()).isNull(); + assertThat(idempotenceManager.lastAckedBatchSequence(tb1)).isEqualTo(Optional.of(1)); + } + + /** + * A genuine out-of-order response (no progress since the batch was sent) must still reset the + * writer. + * + *

Scenario: seq0 is acknowledged so {@code lastAcked=0}. seq1 is then sent while {@code + * lastAcked=0} (no in-flight predecessor). seq1 gets an out-of-order response with no + * advancement since it was sent, indicating an unrecoverable regression, so the writer resets. + */ + @Test + void testGenuineOutOfOrderResponseResetsWriter() throws Exception { + IdempotenceManager idempotenceManager = createIdempotenceManager(true); + Sender sender1 = setupWithIdempotenceState(idempotenceManager); + sender1.runOnce(); + assertThat(idempotenceManager.isWriterIdValid()).isTrue(); + + // Establish lastAcked = 0. + CompletableFuture future0 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(1, "a"), (tb, leo, e) -> future0.complete(e)); + sender1.runOnce(); + finishIdempotentProduceLogRequest(0, tb1, 0, createProduceLogResponse(tb1, 0L, 1L)); + sender1.runOnce(); + assertThat(future0.get()).isNull(); + assertThat(idempotenceManager.lastAckedBatchSequence(tb1)).isEqualTo(Optional.of(0)); + + // Send seq1 while lastAcked is already 0 (no in-flight predecessor). + CompletableFuture future1 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(2, "b"), (tb, leo, e) -> future1.complete(e)); + sender1.runOnce(); + + // seq1 gets an out-of-order response with no progress since it was sent: genuine + // regression -> the writer is reset. Do not run the sender again before asserting, since + // the next iteration would re-initialize a fresh writer id. + finishIdempotentProduceLogRequest( + 1, tb1, 0, createProduceLogResponse(tb1, Errors.OUT_OF_ORDER_SEQUENCE_EXCEPTION)); + assertThat(future1.get()).isInstanceOf(OutOfOrderSequenceException.class); + assertThat(idempotenceManager.isWriterIdValid()).isFalse(); + } + + /** + * The send-time acked-sequence snapshot must be refreshed on every send attempt: after a stale + * out-of-order response is retried, a subsequent out-of-order response with no new + * acknowledgement in between must reset the writer (it is no longer stale). This guards against + * regressing the snapshot to a capture-once semantic. + * + *

Scenario: seq0 and seq1 are in flight (each sent at {@code lastAcked=-1}). seq0 succeeds + * ({@code lastAcked=0}). seq1's first (stale) out-of-order response is retried and the resend + * refreshes seq1's snapshot to 0. A second out-of-order response for seq1, with no new ack, + * then finds {@code 0 > 0} false and resets the writer. + */ + @Test + void testSendTimeSnapshotIsRefreshedOnRetrySoRepeatedOutOfOrderResets() throws Exception { + IdempotenceManager idempotenceManager = createIdempotenceManager(true); + Sender sender1 = setupWithIdempotenceState(idempotenceManager); + sender1.runOnce(); + assertThat(idempotenceManager.isWriterIdValid()).isTrue(); + long writerId = idempotenceManager.writerId(); + + // Send seq0 and seq1; both are sent while lastAcked is not present (-1). + CompletableFuture future0 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(1, "a"), (tb, leo, e) -> future0.complete(e)); + sender1.runOnce(); + CompletableFuture future1 = new CompletableFuture<>(); + appendToAccumulator(tb1, row(2, "b"), (tb, leo, e) -> future1.complete(e)); + sender1.runOnce(); + + // seq0 succeeds first, advancing lastAcked to 0. + finishIdempotentProduceLogRequest(0, tb1, 0, createProduceLogResponse(tb1, 0L, 1L)); + sender1.runOnce(); + assertThat(idempotenceManager.lastAckedBatchSequence(tb1)).isEqualTo(Optional.of(0)); + + // First (stale) seq1 out-of-order: sent at lastAcked=-1 < current 0 -> retried, writer + // kept. The following runOnce resends seq1, refreshing its send-time snapshot to 0. + finishIdempotentProduceLogRequest( + 1, tb1, 0, createProduceLogResponse(tb1, Errors.OUT_OF_ORDER_SEQUENCE_EXCEPTION)); + sender1.runOnce(); + assertThat(idempotenceManager.isWriterIdValid()).isTrue(); + assertThat(idempotenceManager.writerId()).isEqualTo(writerId); + assertThat(future1.isDone()).isFalse(); + + // Second seq1 out-of-order with no new ack: the refreshed snapshot is now 0, so the + // response is no longer stale (0 > 0 is false) and the writer is reset. + finishIdempotentProduceLogRequest( + 1, tb1, 0, createProduceLogResponse(tb1, Errors.OUT_OF_ORDER_SEQUENCE_EXCEPTION)); + assertThat(future1.get()).isInstanceOf(OutOfOrderSequenceException.class); + assertThat(idempotenceManager.isWriterIdValid()).isFalse(); + } + @Test void testCorrectHandlingOfDuplicateSequenceError() throws Exception { IdempotenceManager idempotenceManager = createIdempotenceManager(true);