From cc91031e67903869ba84832729129d81cf1abaca Mon Sep 17 00:00:00 2001 From: fhan Date: Wed, 16 Sep 2026 09:34:26 +0800 Subject: [PATCH 1/6] [client] Back off writes rejected by disk protection When disk protection rejects a write with DISK_WRITE_LOCKED, the Log/KV writers apply a bounded exponential backoff before retrying that bucket, instead of hot-looping and hammering the rejecting server. Write-throttling is unified in a new WriteThrottleController that owns two independent gates per bucket: - KV backpressure (wall-clock, quadratic) for KV upsert/delete - disk-write backoff (monotonic nanos, exponential, never-shorten, ceil-to-ms) for DISK_WRITE_LOCKED rejections Per bucket the effective delay is max(kvRemainingMs, diskRemainingMs); the accumulator wakes on the min across buckets. bucketReady keeps leader-first ordering: a leaderless bucket is still reported to unknownLeaderTables while a gate is pending, so metadata refresh (which sends no data) stays orthogonal to throttling. Adds ConfigOptions for the backoff bounds, documentation, unit coverage in RecordAccumulatorTest/SenderTest, and DiskWriteBackoffITCase. --- .../fluss/client/write/RecordAccumulator.java | 237 +++++------ .../org/apache/fluss/client/write/Sender.java | 29 +- .../client/write/WriteThrottleController.java | 265 ++++++++++++ .../client/table/DiskWriteBackoffITCase.java | 217 ++++++++++ .../client/write/RecordAccumulatorTest.java | 216 ++++++++++ .../apache/fluss/client/write/SenderTest.java | 377 +++++++++++++++++- .../apache/fluss/config/ConfigOptions.java | 22 + tools/maven/suppressions.xml | 1 + website/docs/maintenance/configuration.md | 30 ++ 9 files changed, 1263 insertions(+), 131 deletions(-) create mode 100644 fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java create mode 100644 fluss-client/src/test/java/org/apache/fluss/client/table/DiskWriteBackoffITCase.java diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java index 45f45d89e9..2287920dd1 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java @@ -42,6 +42,7 @@ import org.apache.fluss.shaded.arrow.org.apache.arrow.memory.BufferAllocator; import org.apache.fluss.shaded.arrow.org.apache.arrow.memory.ChunkedAllocationManager; import org.apache.fluss.utils.CopyOnWriteMap; +import org.apache.fluss.utils.ExponentialBackoff; import org.apache.fluss.utils.MathUtils; import org.apache.fluss.utils.clock.Clock; @@ -52,6 +53,7 @@ import javax.annotation.concurrent.GuardedBy; import java.io.IOException; +import java.time.Duration; import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Collections; @@ -64,12 +66,14 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import static org.apache.fluss.record.LogRecordBatchFormat.NO_BATCH_SEQUENCE; import static org.apache.fluss.record.LogRecordBatchFormat.NO_WRITER_ID; import static org.apache.fluss.shaded.arrow.org.apache.arrow.memory.BufferAllocatorUtil.createBufferAllocator; import static org.apache.fluss.utils.PartitionUtils.HISTORICAL_PARTITION_VALUE; +import static org.apache.fluss.utils.Preconditions.checkArgument; import static org.apache.fluss.utils.Preconditions.checkNotNull; /* This file is based on source code of Apache Kafka Project (https://kafka.apache.org/), licensed by the Apache @@ -113,7 +117,7 @@ public final class RecordAccumulator { private final Object resourcesLock = new Object(); @GuardedBy("resourcesLock") - private boolean resourcesDestroyed; + private final AtomicBoolean resourcesDestroyed = new AtomicBoolean(false); /** The pool of lazily created arrow {@link ArrowWriter}s for arrow log write batch. */ private final ArrowWriterPool arrowWriterPool; @@ -133,17 +137,10 @@ public final class RecordAccumulator { private final Clock clock; private final DynamicWriteBatchSizeEstimator batchSizeEstimator; - // Per-bucket backpressure throttle expiry timestamp. Accessed strictly by key on - // hot paths (get / put / remove); writes happen on every backpressure signal and - // every eviction, so the container is sized for lock-striped O(1) updates without - // any whole-map snapshot cost. - private final ConcurrentMap throttleExpiryMs = new ConcurrentHashMap<>(); - private final long maxThrottleMs; - - // Latest Cluster snapshot fed to the metadata-driven throttle sweep. Identity - // equality against this reference short-circuits the sweep when metadata hasn't - // changed. - private volatile Cluster lastClusterRef = Cluster.empty(); + // Unifies the per-bucket write-throttling gates (KV backpressure and disk-write backoff): + // both paths that decide sendability, ready() and drain(), consult it for the effective + // remaining delay so no single gate can be dropped. + private final WriteThrottleController throttle; // TODO add retryBackoffMs to retry the produce request upon receiving an error. // TODO add deliveryTimeoutMs to report success or failure on record delivery. @@ -159,6 +156,23 @@ public final class RecordAccumulator { idempotenceManager, writerMetricGroup, clock, + createDiskWriteBackoff(conf), + BucketAssignerFactory.defaultFactory(conf)); + } + + @VisibleForTesting + RecordAccumulator( + Configuration conf, + IdempotenceManager idempotenceManager, + WriterMetricGroup writerMetricGroup, + Clock clock, + ExponentialBackoff diskWriteBackoff) { + this( + conf, + idempotenceManager, + writerMetricGroup, + clock, + diskWriteBackoff, BucketAssignerFactory.defaultFactory(conf)); } @@ -169,6 +183,23 @@ public final class RecordAccumulator { WriterMetricGroup writerMetricGroup, Clock clock, BucketAssignerFactory bucketAssignerFactory) { + this( + conf, + idempotenceManager, + writerMetricGroup, + clock, + createDiskWriteBackoff(conf), + bucketAssignerFactory); + } + + @VisibleForTesting + RecordAccumulator( + Configuration conf, + IdempotenceManager idempotenceManager, + WriterMetricGroup writerMetricGroup, + Clock clock, + ExponentialBackoff diskWriteBackoff, + BucketAssignerFactory bucketAssignerFactory) { this.bucketAssignerFactory = checkNotNull(bucketAssignerFactory); this.closed = false; this.flushesInProgress = new AtomicInteger(0); @@ -193,11 +224,30 @@ public final class RecordAccumulator { (int) conf.get(ConfigOptions.CLIENT_WRITER_BUFFER_PAGE_SIZE).getBytes()); this.idempotenceManager = idempotenceManager; this.clock = clock; - this.maxThrottleMs = - conf.get(ConfigOptions.CLIENT_WRITER_KV_BACKPRESSURE_MAX_THROTTLE).toMillis(); + this.throttle = + new WriteThrottleController( + conf.get(ConfigOptions.CLIENT_WRITER_KV_BACKPRESSURE_MAX_THROTTLE) + .toMillis(), + checkNotNull(diskWriteBackoff), + clock, + resourcesDestroyed); registerMetrics(writerMetricGroup); } + private static ExponentialBackoff createDiskWriteBackoff(Configuration conf) { + Duration initial = conf.get(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF); + Duration maximum = conf.get(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX); + checkArgument( + initial.compareTo(Duration.ofMillis(1)) >= 0 + && initial.compareTo(maximum) <= 0 + && maximum.compareTo(Duration.ofMillis(Integer.MAX_VALUE)) <= 0, + "%s and %s must satisfy 1ms <= initial <= maximum <= %sms.", + ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF.key(), + ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX.key(), + Integer.MAX_VALUE); + return new ExponentialBackoff(initial.toMillis(), 2, maximum.toMillis(), 0.2); + } + private void registerMetrics(WriterMetricGroup writerMetricGroup) { // memory segment pool related metrics. writerMetricGroup.gauge(MetricNames.WRITER_BUFFER_TOTAL_BYTES, writerBufferPool::totalSize); @@ -306,7 +356,7 @@ private BucketAndWriteBatches createBucketAndWriteBatches( */ public ReadyCheckResult ready(Cluster cluster) { Set readyNodes = new HashSet<>(); - long nextReadyCheckDelayMs = batchTimeoutMs; + long nextReadyCheckDelayMs = Long.MAX_VALUE; Set unknownLeaderTables = new HashSet<>(); // Go table by table so that we can get queue sizes for buckets in a table and calculate // cumulative frequency table (used in bucket assigner). @@ -321,7 +371,14 @@ public ReadyCheckResult ready(Cluster cluster) { nextReadyCheckDelayMs); } - // TODO and the earliest time at which any non-send-able bucket will be ready; + // When all queued buckets are backing off, wait for the earliest effective deadline. + // In particular, a zero batch timeout must not turn this wait into a busy loop. + // Keep the normal polling cadence for idle writers and unresolved metadata. + if (!readyNodes.isEmpty() + || !unknownLeaderTables.isEmpty() + || nextReadyCheckDelayMs == Long.MAX_VALUE) { + nextReadyCheckDelayMs = Math.min(nextReadyCheckDelayMs, batchTimeoutMs); + } return new ReadyCheckResult(readyNodes, nextReadyCheckDelayMs, unknownLeaderTables); } @@ -629,7 +686,7 @@ public void awaitFlushCompletion() throws InterruptedException { */ public void deallocate(WriteBatch batch) { synchronized (resourcesLock) { - if (incomplete.removeIfPresent(batch) && !resourcesDestroyed) { + if (incomplete.removeIfPresent(batch) && !resourcesDestroyed.get()) { writerBufferPool.returnAll(batch.pooledMemorySegments()); } } @@ -763,36 +820,35 @@ private long bucketReady( int bucketId = entry.getKey(); TableBucket tableBucket = cluster.getTableBucket(tableId, targetPath, bucketId); - // If this bucket is throttled, don't mark its node as ready. - // Instead, factor the remaining throttle time into the next check delay. - Long throttleExpiry = throttleExpiryMs.get(tableBucket); - if (throttleExpiry != null) { - long now = clock.milliseconds(); - if (now < throttleExpiry) { - nextReadyCheckDelayMs = Math.min(nextReadyCheckDelayMs, throttleExpiry - now); - continue; - } - // Expired — evict here to reclaim entries for buckets whose deque - // has gone empty and won't reach the drain-time throttle check. - throttleExpiryMs.remove(tableBucket); - } - Integer leader = cluster.leaderFor(tableBucket); if (leader == null) { // This is a bucket for which leader is not known, but messages are // available to send. Note that entries are currently not removed from - // batches when deque is empty. + // batches when deque is empty. Reported regardless of any pending throttle gate: + // a metadata refresh is orthogonal to the gate (it sends no data), and refreshing + // now means the leader is ready to receive once the gate clears. unknownLeaderTables.add(targetPath); - } else { - nextReadyCheckDelayMs = - batchReady( - exhausted, - leader, - waitedTimeMs, - full, - readyNodes, - nextReadyCheckDelayMs); + continue; + } + + // A bucket blocked by any throttle gate (KV backpressure and/or disk-write backoff) + // cannot become ready until every gate clears, so wake no earlier than the latest of + // them. remainingDelayMs() already folds the reasons into a single per-bucket deadline + // via max(); across buckets we keep the earliest wake-up via min(). + long gateMs = throttle.remainingDelayMs(tableBucket); + if (gateMs > 0) { + nextReadyCheckDelayMs = Math.min(nextReadyCheckDelayMs, gateMs); + continue; } + + nextReadyCheckDelayMs = + batchReady( + exhausted, + leader, + waitedTimeMs, + full, + readyNodes, + nextReadyCheckDelayMs); } return nextReadyCheckDelayMs; @@ -1257,8 +1313,8 @@ private List drainBatchesForOneNode(Cluster cluster, Integer no } private boolean shouldSkipBucket(WriteBatch first, TableBucket tableBucket) { - // Backpressure throttle check: skip this bucket if still under throttle - if (isThrottled(tableBucket)) { + // Skip this bucket while any write-throttling gate still blocks it. + if (throttle.isGated(tableBucket)) { return true; } if (idempotenceManager.idempotenceEnabled()) { @@ -1298,80 +1354,38 @@ private boolean shouldSkipBucket(WriteBatch first, TableBucket tableBucket) { return false; } - // ---- Backpressure throttle methods ---- + // ---- Write-throttling delegates (see WriteThrottleController) ---- - /** - * Check if a bucket is currently under backpressure throttle. - * - *

Performs lazy eviction: if the throttle has expired, the entry is removed from the map to - * prevent unbounded growth. - * - * @return true if the bucket should be skipped during drain - */ boolean isThrottled(TableBucket tableBucket) { - Long expiry = throttleExpiryMs.get(tableBucket); - if (expiry == null) { - return false; - } - if (clock.milliseconds() < expiry) { - return true; - } - // Expired — evict to prevent map leak - throttleExpiryMs.remove(tableBucket); - return false; + return throttle.isKvThrottled(tableBucket); + } + + /** Installs disk backoff before the batch retry count is increased by re-enqueueing. */ + long backoffAfterDiskWriteLocked(ReadyWriteBatch batch) { + return throttle.backoffAfterDiskWrite(batch.tableBucket(), batch.writeBatch().attempts()); + } + + /** Returns the remaining disk backoff without changing its deadline. */ + long diskWriteBackoffRemainingMs(TableBucket tableBucket) { + return throttle.diskRemainingMs(tableBucket); + } + + /** Reclaims expired entries even when their queues no longer contain any batches. */ + void maybeEvictExpiredDiskWriteBackoffs() { + throttle.maybeEvictExpiredDiskBackoffs(); + } + + @VisibleForTesting + int diskWriteBackoffCount() { + return throttle.diskBackoffCount(); } - /** - * Update the throttle state for a bucket based on the received pressure signal. - * - *

The delay grows quadratically with pressure: {@code delay = maxThrottleMs * p^2}, where - * {@code p ∈ [0, 1)}. This provides meaningful throttling across the full ramp-up window while - * remaining gentle at low pressure. - * - * @param tableBucket the bucket to update - * @param pressure value in {@code [0, 1)} on the wire; {@code 0} means recovered, positive - * values trigger a throttle window. {@code 1.0f} is reserved as the internal hard-rejection - * value (never sent by the server): the Sender passes it when the server rejected the write - * outright, and it installs the full {@link #maxThrottleMs} window directly. - */ void updateThrottle(TableBucket tableBucket, float pressure) { - if (pressure >= 1f) { - // Hard rejection: stall the bucket for the full max throttle window, bypassing the - // quadratic curve to avoid long-to-float rounding. - throttleExpiryMs.put(tableBucket, clock.milliseconds() + maxThrottleMs); - return; - } - if (pressure > 0f) { - long delay = (long) (maxThrottleMs * pressure * pressure); - if (delay > 0) { - throttleExpiryMs.put(tableBucket, clock.milliseconds() + delay); - return; - } - } - // Recovered or below the meaningful resolution: remove throttle. - // Note: in production, recovery relies on the last throttle window expiring naturally - // (server stops sending the pressure field once p reaches 0). This branch exists as - // defensive completeness and is exercised by unit tests. - throttleExpiryMs.remove(tableBucket); + throttle.updateKvPressure(tableBucket, pressure); } - /** - * Evict throttle entries whose buckets no longer exist in the given cluster (leader unknown, - * partition dropped, table dropped). - * - *

Invoked on every Sender loop with the current cluster snapshot. The identity short-circuit - * makes this an O(1) no-op when metadata hasn't changed, so the actual O(N) walk only runs once - * per real metadata refresh. - */ void maybeEvictStaleThrottles(Cluster cluster) { - if (cluster == lastClusterRef) { - return; - } - lastClusterRef = cluster; - if (throttleExpiryMs.isEmpty()) { - return; - } - throttleExpiryMs.keySet().removeIf(tb -> cluster.leaderFor(tb) == null); + throttle.maybeEvictStaleThrottles(cluster); } private int getDrainIndex(int id) { @@ -1586,15 +1600,16 @@ public void close() { @VisibleForTesting public void destroyResources() { synchronized (resourcesLock) { - if (resourcesDestroyed) { + if (resourcesDestroyed.get()) { return; } - resourcesDestroyed = true; + resourcesDestroyed.compareAndSet(false, true); writerBufferPool.close(); arrowWriterPool.close(); bufferAllocator.close(); chunkedFactory.close(); } + throttle.clearDiskBackoffs(); } /** Per table bucket and write batches. */ 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 9f8b1e8ba2..060c2c8160 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 @@ -247,6 +247,7 @@ private void sendWriteData() throws Exception { // Refresh per-bucket throttle entries against the current cluster snapshot, // dropping any whose bucket has disappeared from metadata. accumulator.maybeEvictStaleThrottles(clusterSnapshot); + accumulator.maybeEvictExpiredDiskWriteBackoffs(); // get the list of buckets with data ready to send. ReadyCheckResult readyCheckResult = accumulator.ready(clusterSnapshot); @@ -359,6 +360,7 @@ private void reEnqueueBatch(ReadyWriteBatch readyWriteBatch) { boolean reEnqueued = accumulator.reEnqueue(readyWriteBatch); maybeRemoveFromInflightBatches(readyWriteBatch); + wakeup(); if (reEnqueued) { // metrics for retry record count. writerMetricGroup @@ -746,13 +748,8 @@ private Set handleWriteBatchException( } else if (canRetry(readyWriteBatch, error.error())) { // if batch failed because of retrievable exception, we need to retry send all those // batches. - LOG.warn( - "Get error write response on table bucket {}, retrying ({} attempts left). Error: {}", - readyWriteBatch.tableBucket(), - retries - writeBatch.attempts(), - error.formatErrMsg()); - if (!idempotenceManager.idempotenceEnabled()) { + prepareWriteRetry(readyWriteBatch, error); reEnqueueBatch(readyWriteBatch); } else if (idempotenceManager.hasWriterId(writeBatch.writerId())) { // If idempotence is enabled only retry the request if the current writer id is @@ -761,6 +758,7 @@ private Set handleWriteBatchException( "Retrying batch to table-bucket {}, Batch sequence : {}", readyWriteBatch.tableBucket(), writeBatch.batchSequence()); + prepareWriteRetry(readyWriteBatch, error); reEnqueueBatch(readyWriteBatch); } else { Exception exception = @@ -802,6 +800,25 @@ private Set handleWriteBatchException( return invalidMetadataTables; } + private void prepareWriteRetry(ReadyWriteBatch batch, ApiError error) { + if (error.error() == Errors.DISK_WRITE_LOCKED) { + long backoffMs = accumulator.backoffAfterDiskWriteLocked(batch); + LOG.warn( + "Get error write response on table bucket {}, disk backoff {} ms " + + "({} attempts left). Error: {}", + batch.tableBucket(), + backoffMs, + retries - batch.writeBatch().attempts(), + error.formatErrMsg()); + } else { + LOG.warn( + "Get error write response on table bucket {}, retrying ({} attempts left). Error: {}", + batch.tableBucket(), + retries - batch.writeBatch().attempts(), + error.formatErrMsg()); + } + } + /** * Rechecks unknown-leader partitions after a bulk metadata update reports {@link * PartitionNotExistException}, and handles missing partitions for tables with historical diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java new file mode 100644 index 0000000000..1b44887981 --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java @@ -0,0 +1,265 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.client.write; + +import org.apache.fluss.annotation.Internal; +import org.apache.fluss.annotation.VisibleForTesting; +import org.apache.fluss.cluster.Cluster; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.utils.ExponentialBackoff; +import org.apache.fluss.utils.clock.Clock; + +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.apache.fluss.utils.Preconditions.checkNotNull; + +/** + * Unifies the write-throttling gates that can delay a bucket from being sent: KV backpressure and + * disk-write backoff. Both gates are per-{@link TableBucket}. + * + *

The two reasons are stored and installed separately because their semantics genuinely differ: + * + *

    + *
  • KV backpressure uses wall-clock deadlines, is latest-wins (a fresher pressure + * signal may shorten or clear the window), and derives its delay quadratically from the + * pressure value. + *
  • Disk-write backoff uses monotonic deadlines (immune to wall-clock shifts), is + * never-shortened under concurrency, and derives its delay from an exponential backoff + * keyed on the batch retry count. + *
+ * + *

What is unified is the read side: a bucket cannot become ready until every active gate + * has cleared, so {@link #remainingDelayMs(TableBucket)} returns the maximum remaining delay across + * reasons and is the single source consulted by both {@code ready()} and {@code drain()}. Adding a + * future throttling reason should only touch this class, not those two paths. + */ +@Internal +final class WriteThrottleController { + + // KV backpressure: wall-clock expiry, latest-wins. Accessed strictly by key on hot paths + // (get / put / remove); the container is sized for lock-striped O(1) updates without any + // whole-map snapshot cost. + private final ConcurrentMap kvThrottleExpiryMs = new ConcurrentHashMap<>(); + private final long maxThrottleMs; + + // Disk protection is independent of KV pressure, whose responses may shorten or clear a + // throttle. Deadlines use monotonic time and are shared by all queues targeting a bucket. + private final ConcurrentMap diskBackoffDeadlineNanos = + new ConcurrentHashMap<>(); + private final ExponentialBackoff diskBackoff; + // Only the sender thread performs periodic sweeps. + private long lastDiskSweepNanos; + + // Latest Cluster snapshot fed to the metadata-driven throttle sweep. Identity equality against + // this reference short-circuits the sweep when metadata hasn't changed. + private volatile Cluster lastClusterRef = Cluster.empty(); + + private final Clock clock; + // Shared with the owning accumulator: a late RPC callback must not retain state after final + // resource destruction. + private final AtomicBoolean resourcesDestroyed; + + WriteThrottleController( + long maxThrottleMs, + ExponentialBackoff diskBackoff, + Clock clock, + AtomicBoolean resourcesDestroyed) { + this.maxThrottleMs = maxThrottleMs; + this.diskBackoff = checkNotNull(diskBackoff); + this.clock = clock; + this.resourcesDestroyed = resourcesDestroyed; + this.lastDiskSweepNanos = clock.nanoseconds(); + } + + // ------------------------------------------------------------------------ + // Unified read side: the single home for "how long until this bucket may send". + // ------------------------------------------------------------------------ + + /** + * Remaining delay before the bucket may be sent, i.e. the latest deadline across all active + * gates. Both reasons express their remainder in milliseconds from now, so the maximum is well + * defined even though they track different clocks internally. Expired entries are evicted + * lazily as a side effect. + */ + long remainingDelayMs(TableBucket tableBucket) { + return Math.max(kvRemainingMs(tableBucket), diskRemainingMs(tableBucket)); + } + + /** Whether any gate currently blocks the bucket from being sent. */ + boolean isGated(TableBucket tableBucket) { + return remainingDelayMs(tableBucket) > 0; + } + + // ------------------------------------------------------------------------ + // KV backpressure (wall-clock millis, latest-wins, quadratic delay). + // ------------------------------------------------------------------------ + + /** + * Update the throttle state for a bucket based on the received pressure signal. + * + *

The delay grows quadratically with pressure: {@code delay = maxThrottleMs * p^2}, where + * {@code p ∈ [0, 1)}. This provides meaningful throttling across the full ramp-up window while + * remaining gentle at low pressure. + * + * @param tableBucket the bucket to update + * @param pressure value in {@code [0, 1)} on the wire; {@code 0} means recovered, positive + * values trigger a throttle window. {@code 1.0f} is reserved as the internal hard-rejection + * value (never sent by the server): the Sender passes it when the server rejected the write + * outright, and it installs the full {@link #maxThrottleMs} window directly. + */ + void updateKvPressure(TableBucket tableBucket, float pressure) { + if (pressure >= 1f) { + // Hard rejection: stall the bucket for the full max throttle window, bypassing the + // quadratic curve to avoid long-to-float rounding. + kvThrottleExpiryMs.put(tableBucket, clock.milliseconds() + maxThrottleMs); + return; + } + if (pressure > 0f) { + long delay = (long) (maxThrottleMs * pressure * pressure); + if (delay > 0) { + kvThrottleExpiryMs.put(tableBucket, clock.milliseconds() + delay); + return; + } + } + // Recovered or below the meaningful resolution: remove throttle. + // Note: in production, recovery relies on the last throttle window expiring naturally + // (server stops sending the pressure field once p reaches 0). This branch exists as + // defensive completeness and is exercised by unit tests. + kvThrottleExpiryMs.remove(tableBucket); + } + + boolean isKvThrottled(TableBucket tableBucket) { + return kvRemainingMs(tableBucket) > 0; + } + + private long kvRemainingMs(TableBucket tableBucket) { + Long expiry = kvThrottleExpiryMs.get(tableBucket); + if (expiry == null) { + return 0; + } + long remainingMs = expiry - clock.milliseconds(); + if (remainingMs > 0) { + return remainingMs; + } + // Expired — evict to prevent map leak. + kvThrottleExpiryMs.remove(tableBucket, expiry); + return 0; + } + + /** + * Evict throttle entries whose buckets no longer exist in the given cluster (leader unknown, + * partition dropped, table dropped). + * + *

Invoked on every Sender loop with the current cluster snapshot. The identity short-circuit + * makes this an O(1) no-op when metadata hasn't changed, so the actual O(N) walk only runs once + * per real metadata refresh. + */ + void maybeEvictStaleThrottles(Cluster cluster) { + if (cluster == lastClusterRef) { + return; + } + lastClusterRef = cluster; + if (kvThrottleExpiryMs.isEmpty()) { + return; + } + kvThrottleExpiryMs.keySet().removeIf(tb -> cluster.leaderFor(tb) == null); + } + + // ------------------------------------------------------------------------ + // Disk protection (monotonic nanos, never-shorten, exponential backoff). + // ------------------------------------------------------------------------ + + /** + * Installs disk backoff for the bucket before the batch retry count is increased by + * re-enqueueing. Concurrent installs never shorten an existing deadline. + * + * @param tableBucket the bucket that was rejected by disk protection + * @param attempts the batch retry count used to derive the exponential backoff + * @return the effective remaining backoff in milliseconds + */ + long backoffAfterDiskWrite(TableBucket tableBucket, int attempts) { + if (resourcesDestroyed.get()) { + return 0; + } + long now = clock.nanoseconds(); + long delayNanos = + TimeUnit.MILLISECONDS.toNanos(Math.max(1L, diskBackoff.backoff(attempts))); + Long deadline = + diskBackoffDeadlineNanos.compute( + tableBucket, + (bucket, previous) -> + previous != null && previous - now > delayNanos + ? previous + : now + delayNanos); + // A late RPC callback must not retain state after final resource destruction. + if (resourcesDestroyed.get()) { + diskBackoffDeadlineNanos.remove(tableBucket, deadline); + return 0; + } + return nanosToCeilMillis(deadline - now); + } + + /** Returns the remaining disk backoff without changing its deadline. */ + long diskRemainingMs(TableBucket tableBucket) { + Long deadline = diskBackoffDeadlineNanos.get(tableBucket); + if (deadline == null) { + return 0; + } + long remainingNanos = deadline - clock.nanoseconds(); + if (remainingNanos > 0) { + return nanosToCeilMillis(remainingNanos); + } + diskBackoffDeadlineNanos.remove(tableBucket, deadline); + return 0; + } + + /** Reclaims expired entries even when their queues no longer contain any batches. */ + void maybeEvictExpiredDiskBackoffs() { + if (diskBackoffDeadlineNanos.isEmpty()) { + return; + } + long now = clock.nanoseconds(); + if (now - lastDiskSweepNanos < TimeUnit.SECONDS.toNanos(1)) { + return; + } + lastDiskSweepNanos = now; + diskBackoffDeadlineNanos.forEach( + (bucket, deadline) -> { + if (deadline - now <= 0) { + diskBackoffDeadlineNanos.remove(bucket, deadline); + } + }); + } + + /** Drops all disk-backoff state; mirrors the accumulator's resource destruction. */ + void clearDiskBackoffs() { + diskBackoffDeadlineNanos.clear(); + } + + @VisibleForTesting + int diskBackoffCount() { + return diskBackoffDeadlineNanos.size(); + } + + private static long nanosToCeilMillis(long nanos) { + return 1 + (nanos - 1) / TimeUnit.MILLISECONDS.toNanos(1); + } +} diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/DiskWriteBackoffITCase.java b/fluss-client/src/test/java/org/apache/fluss/client/table/DiskWriteBackoffITCase.java new file mode 100644 index 0000000000..8870878e9b --- /dev/null +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/DiskWriteBackoffITCase.java @@ -0,0 +1,217 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.client.table; + +import org.apache.fluss.client.Connection; +import org.apache.fluss.client.ConnectionFactory; +import org.apache.fluss.client.FlussConnection; +import org.apache.fluss.client.admin.Admin; +import org.apache.fluss.client.metrics.TestingWriterMetricGroup; +import org.apache.fluss.client.table.scanner.log.LogScanner; +import org.apache.fluss.client.write.WriteFormat; +import org.apache.fluss.client.write.WriteRecord; +import org.apache.fluss.client.write.WriterClient; +import org.apache.fluss.config.ConfigOptions; +import org.apache.fluss.config.Configuration; +import org.apache.fluss.config.MemorySize; +import org.apache.fluss.metadata.DatabaseDescriptor; +import org.apache.fluss.metadata.PhysicalTablePath; +import org.apache.fluss.metadata.TableDescriptor; +import org.apache.fluss.metadata.TableInfo; +import org.apache.fluss.metadata.TablePath; +import org.apache.fluss.row.BinaryRow; +import org.apache.fluss.row.encode.CompactedKeyEncoder; +import org.apache.fluss.server.replica.ReplicaManager; +import org.apache.fluss.server.testutils.FlussClusterExtension; + +import org.junit.jupiter.api.extension.RegisterExtension; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.time.Duration; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +import static org.apache.fluss.record.TestData.DATA1_ROW_TYPE; +import static org.apache.fluss.record.TestData.DATA1_SCHEMA; +import static org.apache.fluss.record.TestData.DATA1_SCHEMA_PK; +import static org.apache.fluss.testutils.DataTestUtils.compactedRow; +import static org.apache.fluss.testutils.DataTestUtils.row; +import static org.apache.fluss.testutils.InternalRowAssert.assertThatRow; +import static org.apache.fluss.testutils.common.CommonTestUtils.retry; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Real RPC tests for retrying writes rejected by disk protection. */ +class DiskWriteBackoffITCase { + + @RegisterExtension + static final FlussClusterExtension CLUSTER = + FlussClusterExtension.builder() + .setNumOfTabletServers(1) + .setClusterConf(clusterConfig()) + .build(); + + @ParameterizedTest + @CsvSource({"false,false", "true,false", "false,true"}) + void testDiskProtectionBackoffAndRecovery(boolean kv, boolean closeWhileLocked) + throws Exception { + Configuration conf = CLUSTER.getClientConfig(); + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_TIMEOUT, Duration.ZERO); + conf.set(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF, Duration.ofSeconds(1)); + conf.set(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX, Duration.ofSeconds(1)); + TestingWriterMetricGroup metrics = TestingWriterMetricGroup.newInstance(); + ReplicaManager replicas = CLUSTER.getTabletServers().iterator().next().getReplicaManager(); + TablePath path = + TablePath.of( + "disk_backoff", + closeWhileLocked ? "close_table" : kv ? "kv_table" : "log_table"); + try (Connection connection = ConnectionFactory.createConnection(conf); + Admin admin = connection.getAdmin()) { + admin.createDatabase(path.getDatabaseName(), DatabaseDescriptor.EMPTY, true).get(); + admin.createTable( + path, + TableDescriptor.builder() + .schema(kv ? DATA1_SCHEMA_PK : DATA1_SCHEMA) + .distributedBy(1) + .build(), + false) + .get(); + try (Table table = connection.getTable(path)) { + WriterClient writer = + new WriterClient( + conf, + ((FlussConnection) connection).getMetadataUpdater(), + metrics, + admin); + try { + // Establish a writable leader before simulating disk protection. + send(writer, table.getTableInfo(), kv, 1).get(30, TimeUnit.SECONDS); + long retriesBeforeRejection = metrics.recordsRetryTotal().getCount(); + replicas.getDiskUsageMonitor().updateWriteLimitConfig(0.85, 0.80); + replicas.getDiskUsageMonitor().update(0.99); + assertThat(replicas.isDiskWriteLocked()).isTrue(); + CompletableFuture rejected = send(writer, table.getTableInfo(), kv, 2); + retry( + Duration.ofSeconds(10), + () -> + assertThat(metrics.recordsRetryTotal().getCount()) + .isEqualTo(retriesBeforeRejection + 1)); + long bytesAfterRejection = metrics.bytesSendTotal().getCount(); + // The callback remains pending, and the payload is not sent again in the + // window. + assertThatThrownBy(() -> rejected.get(200, TimeUnit.MILLISECONDS)) + .isInstanceOf(TimeoutException.class); + assertThat(metrics.recordsRetryTotal().getCount()) + .isEqualTo(retriesBeforeRejection + 1); + assertThat(metrics.bytesSendTotal().getCount()).isEqualTo(bytesAfterRejection); + if (closeWhileLocked) { + long start = System.nanoTime(); + writer.close(Duration.ofMillis(50)); + assertThat(Duration.ofNanos(System.nanoTime() - start)) + .isLessThan(Duration.ofSeconds(2)); + assertThat(metrics.recordsRetryTotal().getCount()) + .isEqualTo(retriesBeforeRejection + 1); + return; + } + replicas.getDiskUsageMonitor().update(0.5); + assertThat(replicas.isDiskWriteLocked()).isFalse(); + rejected.get(10, TimeUnit.SECONDS); + writer.flush(); + if (kv) { + assertThatRow( + table.newLookup() + .createLookuper() + .lookup(row(2)) + .get() + .getSingletonRow()) + .withSchema(DATA1_ROW_TYPE) + .isEqualTo(row(2, "payload")); + } else { + try (LogScanner scanner = table.newScan().createLogScanner()) { + scanner.subscribeFromBeginning(0); + CompletableFuture read = new CompletableFuture<>(); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + while (!read.isDone() && System.nanoTime() < deadline) { + scanner.poll(Duration.ofMillis(100)) + .forEach( + record -> { + if (record.getRow().getInt(0) == 2) { + assertThatRow(record.getRow()) + .withSchema(DATA1_ROW_TYPE) + .isEqualTo(row(2, "payload")); + read.complete(null); + } + }); + } + assertThat(read).isCompleted(); + } + } + } finally { + replicas.getDiskUsageMonitor().update(0.5); + writer.close(Duration.ofSeconds(5)); + } + } + } + } + + private static CompletableFuture send( + WriterClient writer, TableInfo info, boolean kv, int key) { + PhysicalTablePath path = PhysicalTablePath.of(info.getTablePath()); + WriteRecord record; + if (kv) { + BinaryRow value = compactedRow(DATA1_ROW_TYPE, new Object[] {key, "payload"}); + byte[] encodedKey = + new CompactedKeyEncoder(DATA1_ROW_TYPE, DATA1_SCHEMA_PK.getPrimaryKeyIndexes()) + .encodeKey(value); + record = + WriteRecord.forUpsert( + info, + path, + value, + encodedKey, + encodedKey, + WriteFormat.COMPACTED_KV, + null); + } else { + record = WriteRecord.forArrowAppend(info, path, row(key, "payload"), null); + } + CompletableFuture result = new CompletableFuture<>(); + writer.send( + record, + (bucket, offset, error) -> { + if (error == null) { + result.complete(null); + } else { + result.completeExceptionally(error); + } + }); + return result; + } + + private static Configuration clusterConfig() { + Configuration conf = new Configuration(); + conf.set(ConfigOptions.DEFAULT_REPLICATION_FACTOR, 1); + // Keep the real periodic sampler from overwriting the simulated usage during each test. + conf.set(ConfigOptions.SERVER_DATA_DISK_CHECK_INTERVAL, Duration.ofHours(1)); + conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_MEMORY_SIZE, MemorySize.parse("1mb")); + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_SIZE, MemorySize.parse("1kb")); + return conf; + } +} diff --git a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java index 50f42464c8..e47e29d2e7 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java @@ -51,6 +51,7 @@ import org.apache.fluss.rpc.gateway.TabletServerGateway; import org.apache.fluss.rpc.metrics.TestingClientMetricGroup; import org.apache.fluss.utils.CloseableIterator; +import org.apache.fluss.utils.ExponentialBackoff; import org.apache.fluss.utils.clock.ManualClock; import org.apache.commons.lang3.RandomStringUtils; @@ -1019,4 +1020,219 @@ void testThrottledBucketSkippedInDrain() throws Exception { .extracting(ReadyWriteBatch::tableBucket) .containsExactly(tb2); } + + @Test + void testDiskBackoffIncreasesAndExpiresAtDeadline() throws Exception { + RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); + ReadyWriteBatch batch = appendAndDrain(accum, 0); + for (long delay : new long[] {1000, 2000, 4000, 8000, 10000, 10000}) { + assertThat(accum.backoffAfterDiskWriteLocked(batch)).isEqualTo(delay); + accum.reEnqueue(batch); + assertThat(accum.ready(cluster).readyNodes).isEmpty(); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(delay); + // Repeated readiness checks must neither resample nor extend the deadline. + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(delay); + clock.advanceTime(Duration.ofMillis(delay - 1)); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isOne(); + assertThat(accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE)) + .isEmpty(); + clock.advanceTime(Duration.ofMillis(1)); + assertThat(accum.ready(cluster).readyNodes).containsExactly(node1.id()); + assertThat( + accum.drain( + cluster, + Collections.singleton(node1.id()), + Integer.MAX_VALUE) + .get(node1.id())) + .extracting(ReadyWriteBatch::writeBatch) + .containsExactly(batch.writeBatch()); + } + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + + @Test + void testDiskBackoffDoesNotBlockHealthyBucketsOrUnknownMetadata() throws Exception { + RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); + ReadyWriteBatch blocked = appendAndDrain(accum, 0); + accum.backoffAfterDiskWriteLocked(blocked); + accum.reEnqueue(blocked); + appendDiskTestRecord(accum, 1); + appendDiskTestRecord(accum, 2); + assertThat(accum.ready(cluster).readyNodes) + .containsExactlyInAnyOrder(node1.id(), node2.id()); + Map> drained = + accum.drain( + cluster, + new HashSet<>(Arrays.asList(node1.id(), node2.id())), + Integer.MAX_VALUE); + assertThat(drained.get(node1.id())) + .extracting(ReadyWriteBatch::tableBucket) + .containsExactly(tb2); + assertThat(drained.get(node2.id())) + .extracting(ReadyWriteBatch::tableBucket) + .containsExactly(tb3); + // Even a caller supplying a ready node cannot bypass the disk gate. + assertThat(accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE)) + .isEmpty(); + Cluster unknownLeaderCluster = updateCluster(Collections.singletonList(bucket2)); + assertThat(accum.ready(unknownLeaderCluster).unknownLeaderTables) + .contains(DATA1_PHYSICAL_TABLE_PATH); + assertThat(accum.ready(unknownLeaderCluster).nextReadyCheckDelayMs).isZero(); + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + + @Test + void testDiskAndKvBackoffAreIndependentDuringFlushAndClose() throws Exception { + RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); + ReadyWriteBatch batch = appendAndDrain(accum, 0); + accum.backoffAfterDiskWriteLocked(batch); + accum.reEnqueue(batch); + accum.updateThrottle(tb1, 1.0f); // KV default is three seconds. + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(3000); + accum.updateThrottle(tb1, 0.1f); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); + accum.updateThrottle(tb1, 0f); + accum.beginFlush(); + accum.close(); + assertThat(accum.ready(cluster).readyNodes).isEmpty(); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); + assertThat(accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE)) + .isEmpty(); + clock.advanceTime(Duration.ofSeconds(1)); + assertThat(accum.ready(cluster).readyNodes).containsExactly(node1.id()); + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + assertThat(accum.diskWriteBackoffCount()).isZero(); + assertThat(accum.backoffAfterDiskWriteLocked(batch)).isZero(); + } + + @Test + void testConcurrentDiskBackoffsNeverShortenDeadlineAndSweepEmptyQueues() throws Exception { + RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); + ReadyWriteBatch shortBatch = appendAndDrain(accum, 0); + ReadyWriteBatch longBatch = appendAndDrain(accum, 0); + for (int i = 0; i < 4; i++) { + longBatch.writeBatch().reEnqueued(); + } + for (int i = 0; i < 50; i++) { + CompletableFuture.allOf( + CompletableFuture.runAsync( + () -> accum.backoffAfterDiskWriteLocked(shortBatch)), + CompletableFuture.runAsync( + () -> accum.backoffAfterDiskWriteLocked(longBatch)), + CompletableFuture.runAsync(accum::maybeEvictExpiredDiskWriteBackoffs)) + .get(); + assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isEqualTo(10000); + } + clock.advanceTime(Duration.ofSeconds(10)); + // Race expiry cleanup with a new rejection. Conditional removal must preserve the new gate. + CompletableFuture.allOf( + CompletableFuture.runAsync(accum::maybeEvictExpiredDiskWriteBackoffs), + CompletableFuture.runAsync(() -> accum.diskWriteBackoffRemainingMs(tb1)), + CompletableFuture.runAsync( + () -> accum.backoffAfterDiskWriteLocked(longBatch))) + .get(); + assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isEqualTo(10000); + assertThat(accum.hasUnDrained()).isFalse(); + clock.advanceTime(Duration.ofSeconds(10)); + accum.maybeEvictExpiredDiskWriteBackoffs(); + assertThat(accum.diskWriteBackoffCount()).isZero(); + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + + @Test + void testDiskBackoffHandlesNanosecondWrapAndRoundsUp() throws Exception { + clock.advanceTime(Long.MAX_VALUE - clock.nanoseconds() - 500000, TimeUnit.NANOSECONDS); + RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1, 2, 1, 0)); + ReadyWriteBatch batch = appendAndDrain(accum, 0); + assertThat(accum.backoffAfterDiskWriteLocked(batch)).isOne(); + clock.advanceTime(999999, TimeUnit.NANOSECONDS); + assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isOne(); + clock.advanceTime(1, TimeUnit.NANOSECONDS); + assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isZero(); + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + + @Test + void testProductionDiskBackoffJitterAndMinimumDelay() throws Exception { + RecordAccumulator accum = createDiskAccumulator(null); + ReadyWriteBatch batch = appendAndDrain(accum, 0); + for (int attempt = 0; attempt < 8; attempt++) { + long baseDelay = Math.min(1000L << attempt, 10000L); + long delay = accum.backoffAfterDiskWriteLocked(batch); + assertThat(delay) + .isBetween( + (long) (baseDelay * 0.8), Math.min((long) (baseDelay * 1.2), 10000L)); + clock.advanceTime(Duration.ofMillis(delay)); + batch.writeBatch().reEnqueued(); + } + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + + accum = createDiskAccumulator(new ExponentialBackoff(0, 2, 0, 0)); + batch = appendAndDrain(accum, 0); + assertThat(accum.backoffAfterDiskWriteLocked(batch)).isOne(); + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + + @Test + void testInvalidDiskBackoffConfiguration() { + Duration[][] invalid = { + {Duration.ZERO, Duration.ofSeconds(10)}, + {Duration.ofNanos(999999), Duration.ofSeconds(10)}, + {Duration.ofMillis(-1), Duration.ofSeconds(10)}, + {Duration.ofSeconds(11), Duration.ofSeconds(10)}, + {Duration.ofMillis(1), Duration.ofMillis((long) Integer.MAX_VALUE + 1)}, + {Duration.ofMillis(1), Duration.ofSeconds(Long.MAX_VALUE)} + }; + for (Duration[] values : invalid) { + conf.set(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF, values[0]); + conf.set(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX, values[1]); + assertThatThrownBy(() -> createDiskAccumulator(null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("client.writer.disk-write-locked.backoff"); + } + } + + private RecordAccumulator createDiskAccumulator(ExponentialBackoff backoff) { + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_TIMEOUT, Duration.ZERO); + conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_MEMORY_SIZE, new MemorySize(1024 * 1024)); + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_SIZE, new MemorySize(1024)); + conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_PAGE_SIZE, new MemorySize(256)); + IdempotenceManager manager = new IdempotenceManager(false, 5, null, null); + return backoff == null + ? new RecordAccumulator( + conf, + manager, + TestingWriterMetricGroup.newInstance(), + clock, + (tableInfo, path) -> bucketAssigner) + : new RecordAccumulator( + conf, + manager, + TestingWriterMetricGroup.newInstance(), + clock, + backoff, + (tableInfo, path) -> bucketAssigner); + } + + private void appendDiskTestRecord(RecordAccumulator accum, int bucket) throws Exception { + bucketAssigner.setBucketId(bucket); + accum.append( + createRecord(indexedRow(DATA1_ROW_TYPE, new Object[] {1, "a"})), + (tb, offset, error) -> {}, + cluster); + } + + private ReadyWriteBatch appendAndDrain(RecordAccumulator accum, int bucket) throws Exception { + appendDiskTestRecord(accum, bucket); + return accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE) + .get(node1.id()) + .get(0); + } } 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 2270d611df..2ee1815208 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 @@ -27,6 +27,7 @@ import org.apache.fluss.config.Configuration; import org.apache.fluss.config.MemorySize; import org.apache.fluss.exception.AuthorizationException; +import org.apache.fluss.exception.DiskWriteLockedException; import org.apache.fluss.exception.FlussRuntimeException; import org.apache.fluss.exception.InvalidBucketRoutingException; import org.apache.fluss.exception.NetworkException; @@ -35,6 +36,7 @@ import org.apache.fluss.exception.TableNotExistException; import org.apache.fluss.exception.TimeoutException; import org.apache.fluss.exception.TooManyPartitionsException; +import org.apache.fluss.exception.UnknownWriterIdException; import org.apache.fluss.metadata.DataLakeFormat; import org.apache.fluss.metadata.PhysicalTablePath; import org.apache.fluss.metadata.Schema; @@ -59,6 +61,9 @@ import org.apache.fluss.server.entity.PutKvDataForBucket; import org.apache.fluss.server.tablet.TestTabletServerGateway; import org.apache.fluss.types.DataTypes; +import org.apache.fluss.utils.ExponentialBackoff; +import org.apache.fluss.utils.clock.Clock; +import org.apache.fluss.utils.clock.ManualClock; import org.apache.fluss.utils.clock.SystemClock; import org.junit.jupiter.api.AfterEach; @@ -81,6 +86,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -129,6 +135,7 @@ final class SenderTest { private TestingBucketAssigner bucketAssigner; private Sender sender = null; private TestingWriterMetricGroup writerMetricGroup; + private Clock clock = SystemClock.getInstance(); // TODO add more tests as kafka SenderTest. @@ -334,8 +341,12 @@ void testAbortsOnlyMissingPartitionWhenRerouteIsUnsafe() throws Exception { assertThat(activeFuture.get()).isNull(); } - @Test - void testNormalAndHistoricalPutRequests() throws Exception { + @ParameterizedTest + @CsvSource({"false,false", "true,false", "true,true"}) + void testNormalAndHistoricalPutRequests(boolean diskRejected, boolean rpcFailure) + throws Exception { + ManualClock manualClock = new ManualClock(); + clock = manualClock; sender.destroyResources(); TableInfo tableInfo = createHistoricalTableInfo(); PhysicalTablePath activePath = PhysicalTablePath.of(tableInfo.getTablePath(), "20990101"); @@ -394,6 +405,41 @@ void testNormalAndHistoricalPutRequests() throws Exception { secondOriginalPath.getPartitionName()); gateway.response(0, createPutKvResponse(activeBucket, 1L)); + if (diskRejected) { + if (rpcFailure) { + gateway.failRequest(0, new DiskWriteLockedException("disk full")); + } else { + gateway.response( + 0, + makePutKvResponse( + Arrays.asList( + PutKvResultForBucket.historicalFailure( + historicalBucket, + Errors.DISK_WRITE_LOCKED.toApiError(), + firstOriginalPath.getPartitionName()), + PutKvResultForBucket.historicalFailure( + historicalBucket, + Errors.DISK_WRITE_LOCKED.toApiError(), + secondOriginalPath.getPartitionName())))); + } + assertThat(accumulator.diskWriteBackoffRemainingMs(historicalBucket)).isEqualTo(1000); + assertThat(accumulator.diskWriteBackoffRemainingMs(activeBucket)).isZero(); + sender.wakeup(); + sender.runOnce(); + assertThat(gateway.pendingRequestSize()).isZero(); + manualClock.advanceTime(Duration.ofSeconds(1)); + sender.runOnce(); + assertThat(gateway.pendingRequestSize()).isOne(); + PutKvRequest retried = (PutKvRequest) gateway.getRequest(0); + assertThat(retried.getBucketsReqsCount()).isEqualTo(2); + Set retriedOriginals = new HashSet<>(); + for (int i = 0; i < retried.getBucketsReqsCount(); i++) { + retriedOriginals.add(retried.getBucketsReqAt(i).getOriginalPartitionName()); + assertThat(retried.getBucketsReqAt(i).getPartitionId()) + .isEqualTo(historicalBucket.getPartitionId()); + } + assertThat(retriedOriginals).isEqualTo(originalPartitionNames); + } gateway.response( 0, makePutKvResponse( @@ -1626,8 +1672,13 @@ void testPutKvStorageExceptionResponseRetriesInsteadOfFailing() throws Exception assertThat(future.get()).isNull(); } - @Test - void testPutKvRpcFailuresRetryOnlyOwnedTableBatches() throws Exception { + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testPutKvRpcFailuresRetryOnlyOwnedTableBatches(boolean diskRejected) throws Exception { + ManualClock manualClock = new ManualClock(); + clock = manualClock; + sender.destroyResources(); + sender = setupWithIdempotenceState(); TablePath secondTablePath = TablePath.of("test_db_2", "test_pk_table_2"); TableInfo secondTableInfo = TableInfo.of( @@ -1662,7 +1713,12 @@ void testPutKvRpcFailuresRetryOnlyOwnedTableBatches() throws Exception { failRequest( tb1, findRequestIndex(tb1, DATA1_TABLE_ID_PK), - new NetworkException("first table request failed")); + diskRejected + ? new DiskWriteLockedException("disk full") + : new NetworkException("first table request failed")); + assertThat(accumulator.diskWriteBackoffRemainingMs(firstTableBucket)) + .isEqualTo(diskRejected ? 1000L : 0L); + assertThat(accumulator.diskWriteBackoffRemainingMs(secondTableBucket)).isZero(); assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isEqualTo(1L); assertThat(sender.numOfInFlightBatches(firstTableBucket)).isEqualTo(0); assertThat(sender.numOfInFlightBatches(secondTableBucket)).isEqualTo(1); @@ -1677,6 +1733,11 @@ void testPutKvRpcFailuresRetryOnlyOwnedTableBatches() throws Exception { assertThat(sender.numOfInFlightBatches(secondTableBucket)).isEqualTo(0); metadataUpdater.updateCluster(clusterBeforeFailures); + if (diskRejected) { + sender.runOnce(); + assertThat(pendingWriteRequestTableIds(tb1)).containsExactly(DATA2_TABLE_ID); + manualClock.advanceTime(Duration.ofSeconds(1)); + } sender.runOnce(); assertThat(pendingWriteRequestTableIds(tb1)) .containsExactlyInAnyOrder(DATA1_TABLE_ID_PK, DATA2_TABLE_ID); @@ -1693,10 +1754,16 @@ void testPutKvRpcFailuresRetryOnlyOwnedTableBatches() throws Exception { assertThat(secondFuture.get()).isNull(); } - @Test - void testSendWhenTableIdChanges() throws Exception { + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testSendWhenTableIdChanges(boolean diskRejected) throws Exception { CompletableFuture future1 = new CompletableFuture<>(); appendToAccumulator(tb1, row(1, "a"), (tb, leo, e) -> future1.complete(e)); + if (diskRejected) { + sender.runOnce(); + failRequest(tb1, 0, new DiskWriteLockedException("disk full")); + assertThat(accumulator.diskWriteBackoffRemainingMs(tb1)).isPositive(); + } TableInfo newTableInfo = TableInfo.of( DATA1_TABLE_PATH, @@ -1707,6 +1774,7 @@ void testSendWhenTableIdChanges() throws Exception { System.currentTimeMillis(), System.currentTimeMillis()); TableBucket newTableBucket = new TableBucket(newTableInfo.getTableId(), tb1.getBucket()); + assertThat(accumulator.diskWriteBackoffRemainingMs(newTableBucket)).isZero(); metadataUpdater.updateTableInfos(Collections.singletonMap(DATA1_TABLE_PATH, newTableInfo)); sender.runOnce(); @@ -1800,14 +1868,32 @@ private static TableInfo createHistoricalTableInfo() { return createPartitionedTableInfo(AutoPartitionTimeUnit.DAY, 7, true); } + private static TableInfo createHistoricalTableInfo( + AutoPartitionTimeUnit timeUnit, int numToRetain) { + return createHistoricalTableInfo(timeUnit, numToRetain, true); + } + + private static TableInfo createHistoricalTableInfo( + AutoPartitionTimeUnit timeUnit, int numToRetain, boolean primaryKey) { + return createPartitionedTableInfo(timeUnit, numToRetain, true, primaryKey); + } + private static TableInfo createPartitionedTableInfo( AutoPartitionTimeUnit timeUnit, int numToRetain, boolean historicalPartitionEnabled) { - Schema schema = - Schema.newBuilder() - .column("id", DataTypes.INT()) - .column("dt", DataTypes.STRING()) - .primaryKey("id", "dt") - .build(); + return createPartitionedTableInfo(timeUnit, numToRetain, historicalPartitionEnabled, true); + } + + private static TableInfo createPartitionedTableInfo( + AutoPartitionTimeUnit timeUnit, + int numToRetain, + boolean historicalPartitionEnabled, + boolean primaryKey) { + Schema.Builder schemaBuilder = + Schema.newBuilder().column("id", DataTypes.INT()).column("dt", DataTypes.STRING()); + if (primaryKey) { + schemaBuilder.primaryKey("id", "dt"); + } + Schema schema = schemaBuilder.build(); TableDescriptor descriptor = TableDescriptor.builder() .schema(schema) @@ -2071,6 +2157,268 @@ private PutKvResponse createPutKvResponse(TableBucket tb, Errors error) { Collections.singletonList(new PutKvResultForBucket(tb, error.toApiError()))); } + @ParameterizedTest + @CsvSource({ + "false,false,false", + "false,false,true", + "false,true,false", + "false,true,true", + "true,false,false", + "true,false,true", + "true,true,false", + "true,true,true" + }) + void testDiskWriteLockedRetries(boolean kv, boolean rpcFailure, boolean idempotent) + throws Exception { + sender.destroyResources(); + ManualClock manualClock = new ManualClock(); + clock = manualClock; + IdempotenceManager manager = createIdempotenceManager(idempotent); + manager.setWriterId(42L); + sender = setupWithIdempotenceState(manager, 6, 0); + TableBucket bucket = kv ? new TableBucket(DATA1_TABLE_ID_PK, 0) : tb1; + CompletableFuture result = appendDiskTestRecord(kv, bucket, 1); + sender.runOnce(); + byte[] originalPayload = diskTestPayload(kv, bucket, getRequest(tb1, 0)); + for (long delay : new long[] {1000, 2000, 4000, 8000, 10000, 10000}) { + ApiMessage request = getRequest(tb1, 0); + assertThat(diskTestPayload(kv, bucket, request)).isEqualTo(originalPayload); + if (!kv && idempotent) { + assertThat( + getProduceLogRecords((ProduceLogRequest) request, bucket) + .batchIterator() + .next() + .writerId()) + .isEqualTo(42L); + assertBatchSequenceEquals(bucket, (ProduceLogRequest) request, 0); + } + if (rpcFailure) { + failRequest( + tb1, 0, new CompletionException(new DiskWriteLockedException("disk full"))); + } else { + finishRequest( + tb1, + 0, + kv + ? createPutKvResponse(bucket, Errors.DISK_WRITE_LOCKED) + : createProduceLogResponse(bucket, Errors.DISK_WRITE_LOCKED)); + } + assertThat(result).isNotDone(); + assertThat(accumulator.ready(metadataUpdater.getCluster()).nextReadyCheckDelayMs) + .isEqualTo(delay); + sender.wakeup(); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isZero(); + manualClock.advanceTime(Duration.ofMillis(delay - 1)); + sender.wakeup(); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isZero(); + manualClock.advanceTime(Duration.ofMillis(1)); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isOne(); + } + finishRequest( + tb1, + 0, + kv ? createPutKvResponse(bucket, 1L) : createProduceLogResponse(bucket, 0L, 1L)); + assertThat(result.get()).isNull(); + assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isEqualTo(6L); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testDiskWriteLockedFailsWhenRetryIsNotAllowed(boolean writerIdChanged) throws Exception { + sender.destroyResources(); + clock = new ManualClock(); + IdempotenceManager manager = createIdempotenceManager(writerIdChanged); + manager.setWriterId(42L); + sender = setupWithIdempotenceState(manager, writerIdChanged ? 5 : 0, 0); + CompletableFuture result = appendDiskTestRecord(false, tb1, 1); + sender.runOnce(); + if (writerIdChanged) { + manager.setWriterId(43L); + } + failRequest(tb1, 0, new DiskWriteLockedException("disk full")); + assertThat(result.get()) + .isInstanceOf( + writerIdChanged + ? UnknownWriterIdException.class + : DiskWriteLockedException.class); + assertThat(accumulator.diskWriteBackoffCount()).isZero(); + assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isZero(); + assertThat(accumulator.hasUnDrained()).isFalse(); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testLateKvSuccessDoesNotClearDiskBackoff(boolean idempotent) throws Exception { + sender.destroyResources(); + ManualClock manualClock = new ManualClock(); + clock = manualClock; + IdempotenceManager manager = createIdempotenceManager(idempotent); + manager.setWriterId(42L); + sender = setupWithIdempotenceState(manager); + TableBucket bucket = new TableBucket(DATA1_TABLE_ID_PK, 0); + CompletableFuture first = appendDiskTestRecord(true, bucket, 1); + sender.runOnce(); + CompletableFuture second = appendDiskTestRecord(true, bucket, 2); + sender.runOnce(); + finishRequest(tb1, 1, createPutKvResponse(bucket, Errors.DISK_WRITE_LOCKED)); + finishRequest(tb1, 0, createPutKvResponse(bucket, 1L, 0f)); + assertThat(first.get()).isNull(); + assertThat(second).isNotDone(); + assertThat(accumulator.diskWriteBackoffRemainingMs(bucket)).isEqualTo(1000); + sender.wakeup(); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isZero(); + manualClock.advanceTime(Duration.ofSeconds(1)); + sender.runOnce(); + finishRequest(tb1, 0, createPutKvResponse(bucket, 2L)); + assertThat(second.get()).isNull(); + } + + @ParameterizedTest + @ValueSource(strings = {"append", "retry", "close"}) + void testDiskBackoffWaitCanBeWokenByHealthyWorkOrClose(String cause) throws Exception { + boolean close = cause.equals("close"); + boolean retryResponse = cause.equals("retry"); + sender.destroyResources(); + clock = new ManualClock(); + sender = setupWithIdempotenceState(); + CompletableFuture blocked = appendDiskTestRecord(false, tb1, 1); + sender.runOnce(); + failRequest(tb1, 0, new DiskWriteLockedException("disk full")); + TableBucket healthy = new TableBucket(DATA1_TABLE_ID, 1); + if (retryResponse) { + appendDiskTestRecord(false, healthy, 2); + sender.runOnce(); + } + // Consume the retry wakeup, then enter a real wait in a controlled sender thread. + sender.runOnce(); + CompletableFuture finished = new CompletableFuture<>(); + Thread thread = + new Thread( + () -> { + try { + sender.runOnce(); + finished.complete(null); + } catch (Throwable t) { + finished.completeExceptionally(t); + } + }, + "disk-backoff-wakeup-test"); + thread.start(); + try { + retry( + Duration.ofSeconds(10), + () -> assertThat(thread.getState()).isEqualTo(Thread.State.TIMED_WAITING)); + if (close) { + sender.forceClose(); + } else if (retryResponse) { + failRequest(healthy, 0, new TimeoutException("retry healthy bucket")); + } else { + appendDiskTestRecord(false, healthy, 2); + sender.wakeup(); // WriterClient wakes the sender after appending a new batch. + } + finished.get(10, TimeUnit.SECONDS); + assertThat(blocked).isNotDone(); + assertThat(pendingRequestSize(tb1)).isZero(); + if (!close) { + sender.runOnce(); + assertThat(pendingRequestSize(healthy)).isOne(); + finishRequest(healthy, 0, createProduceLogResponse(healthy, 0L, 1L)); + } + } finally { + sender.wakeup(); + thread.join(10000); + accumulator.abortAllBatches(new RuntimeException("test cleanup")); + } + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testHistoricalLogDiskRejection(boolean rpcFailure) throws Exception { + sender.destroyResources(); + ManualClock manualClock = new ManualClock(); + clock = manualClock; + TableInfo info = createHistoricalTableInfo(AutoPartitionTimeUnit.DAY, 7, false); + PhysicalTablePath original = PhysicalTablePath.of(info.getTablePath(), "20000101"); + PhysicalTablePath historical = + PhysicalTablePath.of(info.getTablePath(), HISTORICAL_PARTITION_VALUE); + TableBucket bucket = new TableBucket(info.getTableId(), 22L, 0); + metadataUpdater = + new TestingMetadataUpdater(Collections.singletonMap(info.getTablePath(), info)); + metadataUpdater.updateCluster( + partitionedCluster(info, Collections.singletonMap(historical, bucket))); + sender = setupWithIdempotenceState(); + accumulator.checkAndCacheHistoricalPartitionEnabled(info); + accumulator.routeWritesTo(info, original, historical, bucket.getPartitionId()); + CompletableFuture result = new CompletableFuture<>(); + bucketAssigner.setBucketId(bucket.getBucket()); + accumulator.append( + WriteRecord.forArrowAppend(info, original, row(1, "20000101"), null), + (tb, offset, error) -> result.complete(error), + metadataUpdater.getCluster()); + sender.runOnce(); + TestTabletServerGateway gateway = node1Gateway(); + if (rpcFailure) { + gateway.failRequest(0, new DiskWriteLockedException("disk full")); + } else { + gateway.response( + 0, + makeProduceLogResponse( + Collections.singletonList( + ProduceLogResultForBucket.historicalFailure( + bucket, + Errors.DISK_WRITE_LOCKED.toApiError(), + original.getPartitionName())))); + } + sender.wakeup(); + sender.runOnce(); + assertThat(gateway.pendingRequestSize()).isZero(); + manualClock.advanceTime(Duration.ofSeconds(1)); + sender.runOnce(); + ProduceLogRequest request = (ProduceLogRequest) gateway.getRequest(0); + assertThat(request.getBucketsReqAt(0).getPartitionId()).isEqualTo(bucket.getPartitionId()); + assertThat(request.getBucketsReqAt(0).getOriginalPartitionName()) + .isEqualTo(original.getPartitionName()); + gateway.response( + 0, + makeProduceLogResponse( + Collections.singletonList( + ProduceLogResultForBucket.historicalSuccess( + bucket, 0L, 1L, original.getPartitionName())))); + assertThat(result.get()).isNull(); + } + + private CompletableFuture appendDiskTestRecord( + boolean kv, TableBucket bucket, int key) throws Exception { + CompletableFuture result = new CompletableFuture<>(); + if (kv) { + appendKvToAccumulator( + bucket, + compactedRow(DATA1_ROW_TYPE, new Object[] {key, "a"}), + (tb, offset, error) -> result.complete(error)); + } else { + appendToAccumulator( + bucket, row(key, "a"), (tb, offset, error) -> result.complete(error)); + } + return result; + } + + private byte[] diskTestPayload(boolean kv, TableBucket bucket, ApiMessage request) { + java.nio.ByteBuffer buffer = + (kv + ? ((PutKvRequest) request).getBucketsReqAt(0).getRecordsSlice() + : ((ProduceLogRequest) request) + .getBucketsReqAt(0) + .getRecordsSlice()) + .nioBuffer(); + byte[] bytes = new byte[buffer.remaining()]; + buffer.duplicate().get(bytes); + return bytes; + } + private Sender setupWithIdempotenceState() { return setupWithIdempotenceState(createIdempotenceManager(false)); } @@ -2092,7 +2440,8 @@ private Sender setupWithIdempotenceState( conf, idempotenceManager, writerMetricGroup, - SystemClock.getInstance(), + clock, + new ExponentialBackoff(1000, 2, 10000, 0), (tableInfo, path) -> bucketAssigner); return new Sender( accumulator, diff --git a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java index e497ea0703..d1725bbd76 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java @@ -1493,6 +1493,28 @@ public class ConfigOptions { "Setting a value greater than zero will cause the client to resend any record whose " + "send fails with a potentially transient error."); + public static final ConfigOption CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF = + key("client.writer.disk-write-locked.backoff") + .durationType() + .defaultValue(Duration.ofSeconds(1)) + .withDescription( + "The initial retry backoff when disk write protection rejects a write. " + + "Applies to both Log and KV writes, independently of KV backpressure. " + + "The backoff doubles with the batch retry count and uses 20% jitter, " + + "up to client.writer.disk-write-locked.backoff-max. " + + "The value must be at least 1ms and no greater than the maximum backoff."); + + public static final ConfigOption CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX = + key("client.writer.disk-write-locked.backoff-max") + .durationType() + .defaultValue(Duration.ofSeconds(10)) + .withDescription( + "The maximum retry backoff for a bucket whose writes are rejected by disk " + + "write protection. Must be between the initial backoff and 2147483647ms. " + + "Setting this equal to the initial backoff uses a fixed delay without jitter. " + + "Writes retry automatically after the delay, subject to client.writer.retries. " + + "Flush and graceful close also respect this delay."); + public static final ConfigOption CLIENT_WRITER_ENABLE_IDEMPOTENCE = key("client.writer.enable-idempotence") .booleanType() diff --git a/tools/maven/suppressions.xml b/tools/maven/suppressions.xml index 8b8b6d5ab0..26398d5039 100644 --- a/tools/maven/suppressions.xml +++ b/tools/maven/suppressions.xml @@ -28,6 +28,7 @@ then remove these FileLength suppressions. Track: https://github.com/apache/fluss/issues/4312 --> + Date: Wed, 16 Sep 2026 15:55:06 +0800 Subject: [PATCH 2/6] fix(client): invalidate metadata before re-enqueuing write retries --- .../org/apache/fluss/client/write/Sender.java | 38 ++++++++++--------- .../apache/fluss/client/write/SenderTest.java | 37 ++++++++++++++++++ 2 files changed, 57 insertions(+), 18 deletions(-) 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 060c2c8160..93aec22a16 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 @@ -49,6 +49,7 @@ import javax.annotation.concurrent.GuardedBy; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -748,6 +749,25 @@ private Set handleWriteBatchException( } else if (canRetry(readyWriteBatch, error.error())) { // if batch failed because of retrievable exception, we need to retry send all those // batches. + if (error.exception() instanceof InvalidMetadataException) { + if (error.exception() instanceof UnknownTableOrBucketException) { + LOG.warn( + "Received unknown table or bucket error in write request on bucket {}. The table-bucket may not exist.", + readyWriteBatch.tableBucket()); + } else { + LOG.warn( + "Received invalid metadata error in write request on bucket {}. " + + "Going to request metadata update.", + readyWriteBatch.tableBucket(), + error.exception()); + } + // Re-enqueuing publishes the retry to the sender and wakes it up. Invalidate the + // actual RPC target first so the retry cannot race ahead using stale metadata. A + // historical batch remains keyed by its original partition path in the + // accumulator, while its RPC is sent to the internal historical partition. + metadataUpdater.invalidPhysicalTableBucketMeta( + Collections.singleton(writeTargetPath)); + } if (!idempotenceManager.idempotenceEnabled()) { prepareWriteRetry(readyWriteBatch, error); reEnqueueBatch(readyWriteBatch); @@ -769,24 +789,6 @@ private Set handleWriteBatchException( writeBatch.writerId(), idempotenceManager.writerId())); failBatch(readyWriteBatch, exception, false); } - - if (error.exception() instanceof InvalidMetadataException) { - if (error.exception() instanceof UnknownTableOrBucketException) { - LOG.warn( - "Received unknown table or bucket error in write request on bucket {}. The table-bucket may not exist.", - readyWriteBatch.tableBucket()); - } else { - LOG.warn( - "Received invalid metadata error in write request on bucket {}. " - + "Going to request metadata update.", - readyWriteBatch.tableBucket(), - error.exception()); - } - // A historical batch remains keyed by its original partition path in the - // accumulator, but its RPC is sent to the internal historical partition. Invalidate - // the actual RPC target so the retry refreshes the historical bucket metadata. - invalidMetadataTables.add(writeTargetPath); - } } else { LOG.warn( "Get error write response on table bucket {}, fail. Error: {}", 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 2ee1815208..4e3a667442 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 @@ -243,6 +243,43 @@ void testReroutesWriteAfterExplicitMissingPartitionResponse(boolean historicalCo assertThat(future.get()).isNull(); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testInvalidatesMetadataBeforePublishingRetry(boolean kv) throws Exception { + sender.destroyResources(); + Map tableInfos = new HashMap<>(); + tableInfos.put(DATA1_TABLE_PATH, DATA1_TABLE_INFO); + tableInfos.put(DATA1_TABLE_PATH_PK, DATA1_TABLE_INFO_PK); + AtomicReference retryQueuedAtInvalidation = new AtomicReference<>(); + metadataUpdater = + new TestingMetadataUpdater(tableInfos) { + @Override + public void invalidPhysicalTableBucketMeta( + Set physicalTablesToInvalid) { + if (!physicalTablesToInvalid.isEmpty()) { + retryQueuedAtInvalidation.set(accumulator.hasUnDrained()); + } + super.invalidPhysicalTableBucketMeta(physicalTablesToInvalid); + } + }; + sender = setupWithIdempotenceState(); + TableBucket bucket = kv ? new TableBucket(DATA1_TABLE_ID_PK, 0) : tb1; + CompletableFuture result = appendDiskTestRecord(kv, bucket, 1); + sender.runOnce(); + + finishRequest( + tb1, + 0, + kv + ? createPutKvResponse(bucket, Errors.NOT_LEADER_OR_FOLLOWER) + : createProduceLogResponse(bucket, Errors.NOT_LEADER_OR_FOLLOWER)); + + assertThat(retryQueuedAtInvalidation.get()).isFalse(); + assertThat(result).isNotDone(); + assertThat(accumulator.hasUnDrained()).isTrue(); + accumulator.abortAllBatches(new RuntimeException("test cleanup")); + } + @ParameterizedTest @CsvSource({"2, 4", "4, 2"}) void testServerValidatesResolvedBucketCount(int initialCount, int updatedCount) From 769efa6df4a2debcc9cd43ed66b4b28efb65f0b3 Mon Sep 17 00:00:00 2001 From: fhan Date: Tue, 29 Sep 2026 08:49:35 +0800 Subject: [PATCH 3/6] fix(client): wake sender when KV pressure shortens throttle deadline --- .../fluss/client/write/RecordAccumulator.java | 4 +- .../org/apache/fluss/client/write/Sender.java | 9 +++- .../client/write/WriteThrottleController.java | 39 +++++++++----- .../client/write/RecordAccumulatorTest.java | 10 ++-- .../apache/fluss/client/write/SenderTest.java | 52 +++++++++++++++++++ 5 files changed, 94 insertions(+), 20 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java index 2287920dd1..91fbf827a3 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java @@ -1380,8 +1380,8 @@ int diskWriteBackoffCount() { return throttle.diskBackoffCount(); } - void updateThrottle(TableBucket tableBucket, float pressure) { - throttle.updateKvPressure(tableBucket, pressure); + boolean updateThrottle(TableBucket tableBucket, float pressure) { + return throttle.updateKvPressure(tableBucket, pressure); } void maybeEvictStaleThrottles(Cluster cluster) { 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 93aec22a16..7dbc6f9fb3 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 @@ -632,6 +632,7 @@ private void handlePutKvResponse( long tableId, Map writeBatchesByKey) { Set invalidMetadataTablesSet = new HashSet<>(); + boolean eligibilityAdvanced = false; for (PbPutKvRespForBucket respForBucket : putKvResponse.getBucketsRespsList()) { TableBucket tb = new TableBucket( @@ -641,7 +642,7 @@ private void handlePutKvResponse( // Update backpressure throttle from pressure signal if (respForBucket.hasPressure()) { - accumulator.updateThrottle(tb, respForBucket.getPressure()); + eligibilityAdvanced |= accumulator.updateThrottle(tb, respForBucket.getPressure()); } ReadyWriteBatch writeBatch = @@ -667,6 +668,12 @@ private void handlePutKvResponse( } } metadataUpdater.invalidPhysicalTableBucketMeta(invalidMetadataTablesSet); + // A fresher latest-wins KV pressure signal may shorten the effective gate. Wake only after + // the whole response is applied so the sender observes completed batches and invalidated + // metadata together with the new deadline. + if (eligibilityAdvanced) { + wakeup(); + } } private void handleWriteRequestException(Throwable t, List writeBatches) { diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java index 1b44887981..1218bf4fea 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java @@ -124,26 +124,37 @@ boolean isGated(TableBucket tableBucket) { * values trigger a throttle window. {@code 1.0f} is reserved as the internal hard-rejection * value (never sent by the server): the Sender passes it when the server rejected the write * outright, and it installs the full {@link #maxThrottleMs} window directly. + * @return whether this update moved the effective eligibility time earlier, after accounting + * for both the KV and disk gates */ - void updateKvPressure(TableBucket tableBucket, float pressure) { + boolean updateKvPressure(TableBucket tableBucket, float pressure) { + long nowMs = clock.milliseconds(); + long newKvRemainingMs = 0; if (pressure >= 1f) { // Hard rejection: stall the bucket for the full max throttle window, bypassing the // quadratic curve to avoid long-to-float rounding. - kvThrottleExpiryMs.put(tableBucket, clock.milliseconds() + maxThrottleMs); - return; + newKvRemainingMs = maxThrottleMs; + } else if (pressure > 0f) { + newKvRemainingMs = (long) (maxThrottleMs * pressure * pressure); } - if (pressure > 0f) { - long delay = (long) (maxThrottleMs * pressure * pressure); - if (delay > 0) { - kvThrottleExpiryMs.put(tableBucket, clock.milliseconds() + delay); - return; - } + + Long previousExpiryMs; + if (newKvRemainingMs > 0) { + previousExpiryMs = kvThrottleExpiryMs.put(tableBucket, nowMs + newKvRemainingMs); + } else { + // Recovered or below the meaningful resolution: remove throttle. + // Note: in production, recovery relies on the last throttle window expiring naturally + // (server stops sending the pressure field once p reaches 0). This branch exists as + // defensive completeness and is exercised by unit tests. + previousExpiryMs = kvThrottleExpiryMs.remove(tableBucket); } - // Recovered or below the meaningful resolution: remove throttle. - // Note: in production, recovery relies on the last throttle window expiring naturally - // (server stops sending the pressure field once p reaches 0). This branch exists as - // defensive completeness and is exercised by unit tests. - kvThrottleExpiryMs.remove(tableBucket); + + long previousKvRemainingMs = + previousExpiryMs == null ? 0 : Math.max(0, previousExpiryMs - nowMs); + long diskRemainingMs = diskRemainingMs(tableBucket); + long previousEffectiveRemainingMs = Math.max(previousKvRemainingMs, diskRemainingMs); + long newEffectiveRemainingMs = Math.max(newKvRemainingMs, diskRemainingMs); + return newEffectiveRemainingMs < previousEffectiveRemainingMs; } boolean isKvThrottled(TableBucket tableBucket) { diff --git a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java index e47e29d2e7..933f521c82 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java @@ -1089,11 +1089,15 @@ void testDiskAndKvBackoffAreIndependentDuringFlushAndClose() throws Exception { ReadyWriteBatch batch = appendAndDrain(accum, 0); accum.backoffAfterDiskWriteLocked(batch); accum.reEnqueue(batch); - accum.updateThrottle(tb1, 1.0f); // KV default is three seconds. + // Installing a new gate cannot advance eligibility. + assertThat(accum.updateThrottle(tb1, 1.0f)).isFalse(); assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(3000); - accum.updateThrottle(tb1, 0.1f); + // Shortening KV from three seconds to 30 ms advances the effective gate to the one-second + // disk deadline. + assertThat(accum.updateThrottle(tb1, 0.1f)).isTrue(); assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); - accum.updateThrottle(tb1, 0f); + // Clearing KV does not advance eligibility again while the disk gate still dominates. + assertThat(accum.updateThrottle(tb1, 0f)).isFalse(); accum.beginFlush(); accum.close(); assertThat(accum.ready(cluster).readyNodes).isEmpty(); 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 4e3a667442..e857476e29 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 @@ -1657,6 +1657,58 @@ void testPutKvPressureResponseAppliesBucketThrottle() throws Exception { assertThat(accumulator.isThrottled(tableBucket)).isTrue(); } + @Test + void testLowerKvPressureWakesSenderWaitingOnPreviousDeadline() throws Exception { + TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID_PK, 0); + CompletableFuture first = appendDiskTestRecord(true, tableBucket, 1); + sender.runOnce(); + CompletableFuture second = appendDiskTestRecord(true, tableBucket, 2); + sender.runOnce(); + CompletableFuture queued = appendDiskTestRecord(true, tableBucket, 3); + + // The first response installs a 2430 ms gate (3000 * 0.9^2) while another request remains + // in flight and the third batch waits in the accumulator. + finishRequest(tableBucket, 0, createPutKvResponse(tableBucket, 1L, 0.9f)); + assertThat(first.get()).isNull(); + + CompletableFuture finished = new CompletableFuture<>(); + Thread thread = + new Thread( + () -> { + try { + sender.runOnce(); + finished.complete(null); + } catch (Throwable t) { + finished.completeExceptionally(t); + } + }, + "kv-pressure-wakeup-test"); + thread.start(); + try { + retry( + Duration.ofSeconds(10), + () -> assertThat(thread.getState()).isEqualTo(Thread.State.TIMED_WAITING)); + + // Latest-wins pressure shortens the KV gate to about 30 ms. The response must wake the + // sender from its old 2430 ms wait after completing the corresponding in-flight batch. + finishRequest(tableBucket, 0, createPutKvResponse(tableBucket, 2L, 0.1f)); + assertThat(second.get()).isNull(); + finished.get(1, TimeUnit.SECONDS); + assertThat(queued).isNotDone(); + + // One iteration consumes any remaining 30 ms gate; the next sends the queued batch. + sender.runOnce(); + sender.runOnce(); + assertThat(pendingRequestSize(tableBucket)).isOne(); + finishRequest(tableBucket, 0, createPutKvResponse(tableBucket, 3L)); + assertThat(queued.get()).isNull(); + } finally { + sender.wakeup(); + thread.join(10000); + accumulator.abortAllBatches(new RuntimeException("test cleanup")); + } + } + @Test void testPutKvStorageBackpressureResponseAppliesRetryBackoff() throws Exception { TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID_PK, 0); From dd67dff8b855a5ce9ea0d0d6602c0d5135ef8f7d Mon Sep 17 00:00:00 2001 From: fhan Date: Tue, 29 Sep 2026 21:12:13 +0800 Subject: [PATCH 4/6] fix(client): back off retriable writes per bucket --- .../fluss/client/write/RecordAccumulator.java | 78 ++++++++++-- .../org/apache/fluss/client/write/Sender.java | 16 ++- .../client/write/WriteThrottleController.java | 114 +++++++++++++----- .../client/write/RecordAccumulatorTest.java | 90 ++++++++++++-- .../apache/fluss/client/write/SenderTest.java | 48 +++++++- .../apache/fluss/config/ConfigOptions.java | 19 +++ website/docs/maintenance/configuration.md | 30 ++++- 7 files changed, 330 insertions(+), 65 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java index 91fbf827a3..0775f6e37d 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/RecordAccumulator.java @@ -137,12 +137,11 @@ public final class RecordAccumulator { private final Clock clock; private final DynamicWriteBatchSizeEstimator batchSizeEstimator; - // Unifies the per-bucket write-throttling gates (KV backpressure and disk-write backoff): - // both paths that decide sendability, ready() and drain(), consult it for the effective - // remaining delay so no single gate can be dropped. + // Unifies the per-bucket write-throttling gates (KV backpressure, general retry backoff, and + // disk-write backoff): both paths that decide sendability, ready() and drain(), consult it for + // the effective remaining delay so no single gate can be dropped. private final WriteThrottleController throttle; - // TODO add retryBackoffMs to retry the produce request upon receiving an error. // TODO add deliveryTimeoutMs to report success or failure on record delivery. // TODO add nextBatchExpiryTimeMs @@ -156,6 +155,7 @@ public final class RecordAccumulator { idempotenceManager, writerMetricGroup, clock, + createRetriableWriteBackoff(conf), createDiskWriteBackoff(conf), BucketAssignerFactory.defaultFactory(conf)); } @@ -172,6 +172,7 @@ public final class RecordAccumulator { idempotenceManager, writerMetricGroup, clock, + createRetriableWriteBackoff(conf), diskWriteBackoff, BucketAssignerFactory.defaultFactory(conf)); } @@ -188,6 +189,7 @@ public final class RecordAccumulator { idempotenceManager, writerMetricGroup, clock, + createRetriableWriteBackoff(conf), createDiskWriteBackoff(conf), bucketAssignerFactory); } @@ -200,6 +202,25 @@ public final class RecordAccumulator { Clock clock, ExponentialBackoff diskWriteBackoff, BucketAssignerFactory bucketAssignerFactory) { + this( + conf, + idempotenceManager, + writerMetricGroup, + clock, + createRetriableWriteBackoff(conf), + diskWriteBackoff, + bucketAssignerFactory); + } + + @VisibleForTesting + RecordAccumulator( + Configuration conf, + IdempotenceManager idempotenceManager, + WriterMetricGroup writerMetricGroup, + Clock clock, + ExponentialBackoff retryBackoff, + ExponentialBackoff diskWriteBackoff, + BucketAssignerFactory bucketAssignerFactory) { this.bucketAssignerFactory = checkNotNull(bucketAssignerFactory); this.closed = false; this.flushesInProgress = new AtomicInteger(0); @@ -228,12 +249,27 @@ public final class RecordAccumulator { new WriteThrottleController( conf.get(ConfigOptions.CLIENT_WRITER_KV_BACKPRESSURE_MAX_THROTTLE) .toMillis(), + checkNotNull(retryBackoff), checkNotNull(diskWriteBackoff), clock, resourcesDestroyed); registerMetrics(writerMetricGroup); } + private static ExponentialBackoff createRetriableWriteBackoff(Configuration conf) { + Duration initial = conf.get(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF); + Duration maximum = conf.get(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF_MAX); + checkArgument( + !initial.isNegative() + && initial.compareTo(maximum) <= 0 + && maximum.compareTo(Duration.ofMillis(Integer.MAX_VALUE)) <= 0, + "%s and %s must satisfy 0ms <= initial <= maximum <= %sms.", + ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF.key(), + ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF_MAX.key(), + Integer.MAX_VALUE); + return new ExponentialBackoff(initial.toMillis(), 2, maximum.toMillis(), 0.2); + } + private static ExponentialBackoff createDiskWriteBackoff(Configuration conf) { Duration initial = conf.get(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF); Duration maximum = conf.get(ConfigOptions.CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF_MAX); @@ -831,10 +867,10 @@ private long bucketReady( continue; } - // A bucket blocked by any throttle gate (KV backpressure and/or disk-write backoff) - // cannot become ready until every gate clears, so wake no earlier than the latest of - // them. remainingDelayMs() already folds the reasons into a single per-bucket deadline - // via max(); across buckets we keep the earliest wake-up via min(). + // A bucket blocked by any throttle gate cannot become ready until every gate clears, + // so wake no earlier than the latest of them. remainingDelayMs() already folds the + // reasons into a single per-bucket deadline via max(); across buckets we keep the + // earliest wake-up via min(). long gateMs = throttle.remainingDelayMs(tableBucket); if (gateMs > 0) { nextReadyCheckDelayMs = Math.min(nextReadyCheckDelayMs, gateMs); @@ -1360,6 +1396,19 @@ boolean isThrottled(TableBucket tableBucket) { return throttle.isKvThrottled(tableBucket); } + /** + * Installs general retry backoff before the batch retry count is increased by re-enqueueing. + */ + long backoffAfterRetriableWrite(ReadyWriteBatch batch) { + return throttle.backoffAfterRetriableWrite( + batch.tableBucket(), batch.writeBatch().attempts()); + } + + /** Returns the remaining general retry backoff without changing its deadline. */ + long retriableWriteBackoffRemainingMs(TableBucket tableBucket) { + return throttle.retryBackoffRemainingMs(tableBucket); + } + /** Installs disk backoff before the batch retry count is increased by re-enqueueing. */ long backoffAfterDiskWriteLocked(ReadyWriteBatch batch) { return throttle.backoffAfterDiskWrite(batch.tableBucket(), batch.writeBatch().attempts()); @@ -1370,9 +1419,14 @@ long diskWriteBackoffRemainingMs(TableBucket tableBucket) { return throttle.diskRemainingMs(tableBucket); } - /** Reclaims expired entries even when their queues no longer contain any batches. */ - void maybeEvictExpiredDiskWriteBackoffs() { - throttle.maybeEvictExpiredDiskBackoffs(); + /** Reclaims expired backoff entries even when their queues no longer contain any batches. */ + void maybeEvictExpiredWriteBackoffs() { + throttle.maybeEvictExpiredBackoffs(); + } + + @VisibleForTesting + int retriableWriteBackoffCount() { + return throttle.retryBackoffCount(); } @VisibleForTesting @@ -1609,7 +1663,7 @@ public void destroyResources() { bufferAllocator.close(); chunkedFactory.close(); } - throttle.clearDiskBackoffs(); + throttle.clearBackoffs(); } /** Per table bucket and write batches. */ 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 7dbc6f9fb3..34957b0ab0 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 @@ -248,7 +248,7 @@ private void sendWriteData() throws Exception { // Refresh per-bucket throttle entries against the current cluster snapshot, // dropping any whose bucket has disappeared from metadata. accumulator.maybeEvictStaleThrottles(clusterSnapshot); - accumulator.maybeEvictExpiredDiskWriteBackoffs(); + accumulator.maybeEvictExpiredWriteBackoffs(); // get the list of buckets with data ready to send. ReadyCheckResult readyCheckResult = accumulator.ready(clusterSnapshot); @@ -810,19 +810,23 @@ private Set handleWriteBatchException( } private void prepareWriteRetry(ReadyWriteBatch batch, ApiError error) { + long retryBackoffMs = accumulator.backoffAfterRetriableWrite(batch); if (error.error() == Errors.DISK_WRITE_LOCKED) { - long backoffMs = accumulator.backoffAfterDiskWriteLocked(batch); + long diskBackoffMs = accumulator.backoffAfterDiskWriteLocked(batch); LOG.warn( - "Get error write response on table bucket {}, disk backoff {} ms " - + "({} attempts left). Error: {}", + "Get error write response on table bucket {}, disk backoff {} ms and " + + "retry backoff {} ms ({} attempts left). Error: {}", batch.tableBucket(), - backoffMs, + diskBackoffMs, + retryBackoffMs, retries - batch.writeBatch().attempts(), error.formatErrMsg()); } else { LOG.warn( - "Get error write response on table bucket {}, retrying ({} attempts left). Error: {}", + "Get error write response on table bucket {}, retrying after {} ms " + + "({} attempts left). Error: {}", batch.tableBucket(), + retryBackoffMs, retries - batch.writeBatch().attempts(), error.formatErrMsg()); } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java index 1218bf4fea..3183181ae6 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java @@ -32,15 +32,17 @@ import static org.apache.fluss.utils.Preconditions.checkNotNull; /** - * Unifies the write-throttling gates that can delay a bucket from being sent: KV backpressure and - * disk-write backoff. Both gates are per-{@link TableBucket}. + * Unifies the write-throttling gates that can delay a bucket from being sent: KV backpressure, + * retriable-write backoff, and disk-write backoff. All gates are per-{@link TableBucket}. * - *

The two reasons are stored and installed separately because their semantics genuinely differ: + *

The reasons are stored and installed separately because their semantics genuinely differ: * *

    *
  • KV backpressure uses wall-clock deadlines, is latest-wins (a fresher pressure * signal may shorten or clear the window), and derives its delay quadratically from the * pressure value. + *
  • Retriable-write backoff uses monotonic deadlines, is never-shortened under + * concurrency, and applies an exponential delay to every retriable write failure. *
  • Disk-write backoff uses monotonic deadlines (immune to wall-clock shifts), is * never-shortened under concurrency, and derives its delay from an exponential backoff * keyed on the batch retry count. @@ -60,13 +62,19 @@ final class WriteThrottleController { private final ConcurrentMap kvThrottleExpiryMs = new ConcurrentHashMap<>(); private final long maxThrottleMs; + // General retry pacing is independent of policy-specific gates. It covers transient failures + // such as NOT_ENOUGH_REPLICAS and prevents a re-enqueued batch from being resent immediately. + private final ConcurrentMap retryBackoffDeadlineNanos = + new ConcurrentHashMap<>(); + private final ExponentialBackoff retryBackoff; + // Disk protection is independent of KV pressure, whose responses may shorten or clear a // throttle. Deadlines use monotonic time and are shared by all queues targeting a bucket. private final ConcurrentMap diskBackoffDeadlineNanos = new ConcurrentHashMap<>(); private final ExponentialBackoff diskBackoff; // Only the sender thread performs periodic sweeps. - private long lastDiskSweepNanos; + private long lastBackoffSweepNanos; // Latest Cluster snapshot fed to the metadata-driven throttle sweep. Identity equality against // this reference short-circuits the sweep when metadata hasn't changed. @@ -79,14 +87,16 @@ final class WriteThrottleController { WriteThrottleController( long maxThrottleMs, + ExponentialBackoff retryBackoff, ExponentialBackoff diskBackoff, Clock clock, AtomicBoolean resourcesDestroyed) { this.maxThrottleMs = maxThrottleMs; + this.retryBackoff = checkNotNull(retryBackoff); this.diskBackoff = checkNotNull(diskBackoff); this.clock = clock; this.resourcesDestroyed = resourcesDestroyed; - this.lastDiskSweepNanos = clock.nanoseconds(); + this.lastBackoffSweepNanos = clock.nanoseconds(); } // ------------------------------------------------------------------------ @@ -95,12 +105,14 @@ final class WriteThrottleController { /** * Remaining delay before the bucket may be sent, i.e. the latest deadline across all active - * gates. Both reasons express their remainder in milliseconds from now, so the maximum is well + * gates. All reasons express their remainder in milliseconds from now, so the maximum is well * defined even though they track different clocks internally. Expired entries are evicted * lazily as a side effect. */ long remainingDelayMs(TableBucket tableBucket) { - return Math.max(kvRemainingMs(tableBucket), diskRemainingMs(tableBucket)); + return Math.max( + kvRemainingMs(tableBucket), + Math.max(retryBackoffRemainingMs(tableBucket), diskRemainingMs(tableBucket))); } /** Whether any gate currently blocks the bucket from being sent. */ @@ -125,7 +137,7 @@ boolean isGated(TableBucket tableBucket) { * value (never sent by the server): the Sender passes it when the server rejected the write * outright, and it installs the full {@link #maxThrottleMs} window directly. * @return whether this update moved the effective eligibility time earlier, after accounting - * for both the KV and disk gates + * for the KV, retriable-write, and disk gates */ boolean updateKvPressure(TableBucket tableBucket, float pressure) { long nowMs = clock.milliseconds(); @@ -151,9 +163,10 @@ boolean updateKvPressure(TableBucket tableBucket, float pressure) { long previousKvRemainingMs = previousExpiryMs == null ? 0 : Math.max(0, previousExpiryMs - nowMs); - long diskRemainingMs = diskRemainingMs(tableBucket); - long previousEffectiveRemainingMs = Math.max(previousKvRemainingMs, diskRemainingMs); - long newEffectiveRemainingMs = Math.max(newKvRemainingMs, diskRemainingMs); + long otherRemainingMs = + Math.max(retryBackoffRemainingMs(tableBucket), diskRemainingMs(tableBucket)); + long previousEffectiveRemainingMs = Math.max(previousKvRemainingMs, otherRemainingMs); + long newEffectiveRemainingMs = Math.max(newKvRemainingMs, otherRemainingMs); return newEffectiveRemainingMs < previousEffectiveRemainingMs; } @@ -194,6 +207,28 @@ void maybeEvictStaleThrottles(Cluster cluster) { kvThrottleExpiryMs.keySet().removeIf(tb -> cluster.leaderFor(tb) == null); } + // ------------------------------------------------------------------------ + // General retriable-write backoff (monotonic nanos, never-shorten, exponential). + // ------------------------------------------------------------------------ + + /** + * Installs general retry backoff before the batch retry count is increased by re-enqueueing. + * Concurrent installs never shorten an existing deadline. + * + * @param tableBucket the bucket whose write failed with a retriable error + * @param attempts the batch retry count used to derive the exponential backoff + * @return the effective remaining backoff in milliseconds, or zero when disabled + */ + long backoffAfterRetriableWrite(TableBucket tableBucket, int attempts) { + return installNeverShorterBackoff( + retryBackoffDeadlineNanos, tableBucket, retryBackoff.backoff(attempts)); + } + + /** Returns the remaining general retry backoff without changing its deadline. */ + long retryBackoffRemainingMs(TableBucket tableBucket) { + return remainingBackoffMs(retryBackoffDeadlineNanos, tableBucket); + } + // ------------------------------------------------------------------------ // Disk protection (monotonic nanos, never-shorten, exponential backoff). // ------------------------------------------------------------------------ @@ -207,14 +242,24 @@ void maybeEvictStaleThrottles(Cluster cluster) { * @return the effective remaining backoff in milliseconds */ long backoffAfterDiskWrite(TableBucket tableBucket, int attempts) { - if (resourcesDestroyed.get()) { + return installNeverShorterBackoff( + diskBackoffDeadlineNanos, tableBucket, Math.max(1L, diskBackoff.backoff(attempts))); + } + + /** Returns the remaining disk backoff without changing its deadline. */ + long diskRemainingMs(TableBucket tableBucket) { + return remainingBackoffMs(diskBackoffDeadlineNanos, tableBucket); + } + + private long installNeverShorterBackoff( + ConcurrentMap deadlines, TableBucket tableBucket, long delayMs) { + if (delayMs <= 0 || resourcesDestroyed.get()) { return 0; } long now = clock.nanoseconds(); - long delayNanos = - TimeUnit.MILLISECONDS.toNanos(Math.max(1L, diskBackoff.backoff(attempts))); + long delayNanos = TimeUnit.MILLISECONDS.toNanos(delayMs); Long deadline = - diskBackoffDeadlineNanos.compute( + deadlines.compute( tableBucket, (bucket, previous) -> previous != null && previous - now > delayNanos @@ -222,15 +267,15 @@ long backoffAfterDiskWrite(TableBucket tableBucket, int attempts) { : now + delayNanos); // A late RPC callback must not retain state after final resource destruction. if (resourcesDestroyed.get()) { - diskBackoffDeadlineNanos.remove(tableBucket, deadline); + deadlines.remove(tableBucket, deadline); return 0; } return nanosToCeilMillis(deadline - now); } - /** Returns the remaining disk backoff without changing its deadline. */ - long diskRemainingMs(TableBucket tableBucket) { - Long deadline = diskBackoffDeadlineNanos.get(tableBucket); + private long remainingBackoffMs( + ConcurrentMap deadlines, TableBucket tableBucket) { + Long deadline = deadlines.get(tableBucket); if (deadline == null) { return 0; } @@ -238,33 +283,44 @@ long diskRemainingMs(TableBucket tableBucket) { if (remainingNanos > 0) { return nanosToCeilMillis(remainingNanos); } - diskBackoffDeadlineNanos.remove(tableBucket, deadline); + deadlines.remove(tableBucket, deadline); return 0; } - /** Reclaims expired entries even when their queues no longer contain any batches. */ - void maybeEvictExpiredDiskBackoffs() { - if (diskBackoffDeadlineNanos.isEmpty()) { + /** Reclaims expired backoff entries even when their queues no longer contain any batches. */ + void maybeEvictExpiredBackoffs() { + if (retryBackoffDeadlineNanos.isEmpty() && diskBackoffDeadlineNanos.isEmpty()) { return; } long now = clock.nanoseconds(); - if (now - lastDiskSweepNanos < TimeUnit.SECONDS.toNanos(1)) { + if (now - lastBackoffSweepNanos < TimeUnit.SECONDS.toNanos(1)) { return; } - lastDiskSweepNanos = now; - diskBackoffDeadlineNanos.forEach( + lastBackoffSweepNanos = now; + evictExpiredBackoffs(retryBackoffDeadlineNanos, now); + evictExpiredBackoffs(diskBackoffDeadlineNanos, now); + } + + private static void evictExpiredBackoffs(ConcurrentMap deadlines, long now) { + deadlines.forEach( (bucket, deadline) -> { if (deadline - now <= 0) { - diskBackoffDeadlineNanos.remove(bucket, deadline); + deadlines.remove(bucket, deadline); } }); } - /** Drops all disk-backoff state; mirrors the accumulator's resource destruction. */ - void clearDiskBackoffs() { + /** Drops all backoff state; mirrors the accumulator's resource destruction. */ + void clearBackoffs() { + retryBackoffDeadlineNanos.clear(); diskBackoffDeadlineNanos.clear(); } + @VisibleForTesting + int retryBackoffCount() { + return retryBackoffDeadlineNanos.size(); + } + @VisibleForTesting int diskBackoffCount() { return diskBackoffDeadlineNanos.size(); diff --git a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java index 933f521c82..63462e76d9 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/write/RecordAccumulatorTest.java @@ -1051,6 +1051,34 @@ void testDiskBackoffIncreasesAndExpiresAtDeadline() throws Exception { accum.destroyResources(); } + @Test + void testRetriableWriteBackoffIncreasesAndExpiresAtDeadline() throws Exception { + RecordAccumulator accum = + createBackoffAccumulator( + new ExponentialBackoff(100, 2, 1000, 0), + new ExponentialBackoff(1000, 2, 10000, 0)); + ReadyWriteBatch batch = appendAndDrain(accum, 0); + for (long delay : new long[] {100, 200, 400, 800, 1000, 1000}) { + assertThat(accum.backoffAfterRetriableWrite(batch)).isEqualTo(delay); + accum.reEnqueue(batch); + assertThat(accum.ready(cluster).readyNodes).isEmpty(); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(delay); + assertThat(accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE)) + .isEmpty(); + clock.advanceTime(Duration.ofMillis(delay)); + assertThat( + accum.drain( + cluster, + Collections.singleton(node1.id()), + Integer.MAX_VALUE) + .get(node1.id())) + .extracting(ReadyWriteBatch::writeBatch) + .containsExactly(batch.writeBatch()); + } + accum.abortAllBatches(new RuntimeException("test cleanup")); + accum.destroyResources(); + } + @Test void testDiskBackoffDoesNotBlockHealthyBucketsOrUnknownMetadata() throws Exception { RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); @@ -1084,31 +1112,39 @@ void testDiskBackoffDoesNotBlockHealthyBucketsOrUnknownMetadata() throws Excepti } @Test - void testDiskAndKvBackoffAreIndependentDuringFlushAndClose() throws Exception { - RecordAccumulator accum = createDiskAccumulator(new ExponentialBackoff(1000, 2, 10000, 0)); + void testWriteBackoffsAreIndependentDuringFlushAndClose() throws Exception { + RecordAccumulator accum = + createBackoffAccumulator( + new ExponentialBackoff(2000, 2, 2000, 0), + new ExponentialBackoff(1000, 2, 10000, 0)); ReadyWriteBatch batch = appendAndDrain(accum, 0); + accum.backoffAfterRetriableWrite(batch); accum.backoffAfterDiskWriteLocked(batch); accum.reEnqueue(batch); // Installing a new gate cannot advance eligibility. assertThat(accum.updateThrottle(tb1, 1.0f)).isFalse(); assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(3000); - // Shortening KV from three seconds to 30 ms advances the effective gate to the one-second - // disk deadline. + // Shortening KV from three seconds to 30 ms advances the effective gate to the two-second + // general retry deadline. assertThat(accum.updateThrottle(tb1, 0.1f)).isTrue(); - assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); - // Clearing KV does not advance eligibility again while the disk gate still dominates. + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(2000); + // Clearing KV does not advance eligibility again while the retry gate still dominates. assertThat(accum.updateThrottle(tb1, 0f)).isFalse(); accum.beginFlush(); accum.close(); assertThat(accum.ready(cluster).readyNodes).isEmpty(); - assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(2000); assertThat(accum.drain(cluster, Collections.singleton(node1.id()), Integer.MAX_VALUE)) .isEmpty(); clock.advanceTime(Duration.ofSeconds(1)); + assertThat(accum.ready(cluster).nextReadyCheckDelayMs).isEqualTo(1000); + clock.advanceTime(Duration.ofSeconds(1)); assertThat(accum.ready(cluster).readyNodes).containsExactly(node1.id()); accum.abortAllBatches(new RuntimeException("test cleanup")); accum.destroyResources(); + assertThat(accum.retriableWriteBackoffCount()).isZero(); assertThat(accum.diskWriteBackoffCount()).isZero(); + assertThat(accum.backoffAfterRetriableWrite(batch)).isZero(); assertThat(accum.backoffAfterDiskWriteLocked(batch)).isZero(); } @@ -1126,14 +1162,14 @@ void testConcurrentDiskBackoffsNeverShortenDeadlineAndSweepEmptyQueues() throws () -> accum.backoffAfterDiskWriteLocked(shortBatch)), CompletableFuture.runAsync( () -> accum.backoffAfterDiskWriteLocked(longBatch)), - CompletableFuture.runAsync(accum::maybeEvictExpiredDiskWriteBackoffs)) + CompletableFuture.runAsync(accum::maybeEvictExpiredWriteBackoffs)) .get(); assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isEqualTo(10000); } clock.advanceTime(Duration.ofSeconds(10)); // Race expiry cleanup with a new rejection. Conditional removal must preserve the new gate. CompletableFuture.allOf( - CompletableFuture.runAsync(accum::maybeEvictExpiredDiskWriteBackoffs), + CompletableFuture.runAsync(accum::maybeEvictExpiredWriteBackoffs), CompletableFuture.runAsync(() -> accum.diskWriteBackoffRemainingMs(tb1)), CompletableFuture.runAsync( () -> accum.backoffAfterDiskWriteLocked(longBatch))) @@ -1141,7 +1177,7 @@ void testConcurrentDiskBackoffsNeverShortenDeadlineAndSweepEmptyQueues() throws assertThat(accum.diskWriteBackoffRemainingMs(tb1)).isEqualTo(10000); assertThat(accum.hasUnDrained()).isFalse(); clock.advanceTime(Duration.ofSeconds(10)); - accum.maybeEvictExpiredDiskWriteBackoffs(); + accum.maybeEvictExpiredWriteBackoffs(); assertThat(accum.diskWriteBackoffCount()).isZero(); accum.abortAllBatches(new RuntimeException("test cleanup")); accum.destroyResources(); @@ -1203,6 +1239,23 @@ void testInvalidDiskBackoffConfiguration() { } } + @Test + void testInvalidRetriableWriteBackoffConfiguration() { + Duration[][] invalid = { + {Duration.ofMillis(-1), Duration.ofSeconds(1)}, + {Duration.ofSeconds(2), Duration.ofSeconds(1)}, + {Duration.ZERO, Duration.ofMillis((long) Integer.MAX_VALUE + 1)}, + {Duration.ZERO, Duration.ofSeconds(Long.MAX_VALUE)} + }; + for (Duration[] values : invalid) { + conf.set(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF, values[0]); + conf.set(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF_MAX, values[1]); + assertThatThrownBy(() -> createDiskAccumulator(null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("client.writer.retry-backoff"); + } + } + private RecordAccumulator createDiskAccumulator(ExponentialBackoff backoff) { conf.set(ConfigOptions.CLIENT_WRITER_BATCH_TIMEOUT, Duration.ZERO); conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_MEMORY_SIZE, new MemorySize(1024 * 1024)); @@ -1225,6 +1278,23 @@ private RecordAccumulator createDiskAccumulator(ExponentialBackoff backoff) { (tableInfo, path) -> bucketAssigner); } + private RecordAccumulator createBackoffAccumulator( + ExponentialBackoff retryBackoff, ExponentialBackoff diskBackoff) { + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_TIMEOUT, Duration.ZERO); + conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_MEMORY_SIZE, new MemorySize(1024 * 1024)); + conf.set(ConfigOptions.CLIENT_WRITER_BATCH_SIZE, new MemorySize(1024)); + conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_PAGE_SIZE, new MemorySize(256)); + IdempotenceManager manager = new IdempotenceManager(false, 5, null, null); + return new RecordAccumulator( + conf, + manager, + TestingWriterMetricGroup.newInstance(), + clock, + retryBackoff, + diskBackoff, + (tableInfo, path) -> bucketAssigner); + } + private void appendDiskTestRecord(RecordAccumulator accum, int bucket) throws Exception { bucketAssigner.setBucketId(bucket); accum.append( 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 e857476e29..ad6bde5221 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 @@ -1753,8 +1753,9 @@ void testPutKvStorageExceptionResponseRetriesInsteadOfFailing() throws Exception // The batch is re-enqueued for retry rather than completed with an exception. assertThat(future.isDone()).isFalse(); - // Unlike STORAGE_BACKPRESSURE_EXCEPTION, no throttle is installed, so the retried batch - // is sent out again immediately and completes once the server recovers. + // Generic retry backoff is disabled by this test fixture. Unlike + // STORAGE_BACKPRESSURE_EXCEPTION, no policy-specific gate is installed, so the retried + // batch is sent again immediately and completes once the server recovers. sender.runOnce(); assertThat(sender.numOfInFlightBatches(tableBucket)).isEqualTo(1); finishRequest(tableBucket, 0, createPutKvResponse(tableBucket, 1L)); @@ -2314,6 +2315,39 @@ void testDiskWriteLockedRetries(boolean kv, boolean rpcFailure, boolean idempote assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isEqualTo(6L); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testNotEnoughReplicasUsesRetriableWriteBackoff(boolean idempotent) throws Exception { + sender.destroyResources(); + ManualClock manualClock = new ManualClock(); + clock = manualClock; + IdempotenceManager manager = createIdempotenceManager(idempotent); + manager.setWriterId(42L); + sender = setupWithIdempotenceState(manager, Integer.MAX_VALUE, 0, Duration.ofMillis(100)); + CompletableFuture result = appendDiskTestRecord(false, tb1, 1); + sender.runOnce(); + + finishRequest(tb1, 0, createProduceLogResponse(tb1, Errors.NOT_ENOUGH_REPLICAS_EXCEPTION)); + + assertThat(result).isNotDone(); + assertThat(accumulator.retriableWriteBackoffRemainingMs(tb1)).isEqualTo(100); + assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isOne(); + + // Consume the response wakeup. The re-enqueued batch remains blocked by the retry gate. + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isZero(); + manualClock.advanceTime(Duration.ofMillis(99)); + sender.wakeup(); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isZero(); + + manualClock.advanceTime(Duration.ofMillis(1)); + sender.runOnce(); + assertThat(pendingRequestSize(tb1)).isOne(); + finishRequest(tb1, 0, createProduceLogResponse(tb1, 0L, 1L)); + assertThat(result.get()).isNull(); + } + @ParameterizedTest @ValueSource(booleans = {false, true}) void testDiskWriteLockedFailsWhenRetryIsNotAllowed(boolean writerIdChanged) throws Exception { @@ -2518,11 +2552,21 @@ private Sender setupWithIdempotenceState(IdempotenceManager idempotenceManager) private Sender setupWithIdempotenceState( IdempotenceManager idempotenceManager, int reties, int batchTimeoutMs) { + return setupWithIdempotenceState(idempotenceManager, reties, batchTimeoutMs, Duration.ZERO); + } + + private Sender setupWithIdempotenceState( + IdempotenceManager idempotenceManager, + int reties, + int batchTimeoutMs, + Duration retryBackoff) { Configuration conf = new Configuration(); conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_MEMORY_SIZE, new MemorySize(TOTAL_MEMORY_SIZE)); conf.set(ConfigOptions.CLIENT_WRITER_BATCH_SIZE, new MemorySize(BATCH_SIZE)); conf.set(ConfigOptions.CLIENT_WRITER_BUFFER_PAGE_SIZE, new MemorySize(PAGE_SIZE)); conf.set(ConfigOptions.CLIENT_WRITER_BATCH_TIMEOUT, Duration.ofMillis(batchTimeoutMs)); + conf.set(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF, retryBackoff); + conf.set(ConfigOptions.CLIENT_WRITER_RETRY_BACKOFF_MAX, retryBackoff); bucketAssigner = new TestingBucketAssigner(); accumulator = new RecordAccumulator( diff --git a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java index d1725bbd76..d94079f15c 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java @@ -1493,6 +1493,25 @@ public class ConfigOptions { "Setting a value greater than zero will cause the client to resend any record whose " + "send fails with a potentially transient error."); + public static final ConfigOption CLIENT_WRITER_RETRY_BACKOFF = + key("client.writer.retry-backoff") + .durationType() + .defaultValue(Duration.ofMillis(100)) + .withDescription( + "The initial delay before retrying a write that failed with a retriable error. " + + "The delay doubles with the batch retry count and uses 20% jitter, " + + "up to client.writer.retry-backoff-max. Set to 0 to disable this backoff."); + + public static final ConfigOption CLIENT_WRITER_RETRY_BACKOFF_MAX = + key("client.writer.retry-backoff-max") + .durationType() + .defaultValue(Duration.ofSeconds(1)) + .withDescription( + "The maximum delay before retrying a write that failed with a retriable error. " + + "Must be between the initial backoff and 2147483647ms. " + + "The backoff applies per physical table bucket and is independent of " + + "KV backpressure and disk-write protection."); + public static final ConfigOption CLIENT_WRITER_DISK_WRITE_LOCKED_BACKOFF = key("client.writer.disk-write-locked.backoff") .durationType() diff --git a/website/docs/maintenance/configuration.md b/website/docs/maintenance/configuration.md index 28d60db206..5aca72feaa 100644 --- a/website/docs/maintenance/configuration.md +++ b/website/docs/maintenance/configuration.md @@ -220,6 +220,23 @@ The logging-related environment options (`env.log.dir`, `env.log.level`, `env.lo | kv.scanner.max-per-server | Integer | 200 | The maximum total number of concurrent KV scanner sessions allowed across all buckets on a single tablet server. New scan requests that exceed this limit will be rejected with an error. The default value is 200. | | kv.scanner.max-batch-size | MemorySize | 10mb | Server-side cap on the per-batch payload size for KV full-scan responses. The effective batch size is min(client-requested batch_size_bytes, this value). Protects the tablet server from out-of-memory if a client passes an excessively large batch size. The default value is 10mb. | +## Client writer retry backoff + +| Key | Default | Type | Description | +| :--- | :--- | :--- | :--- | +| `client.writer.retry-backoff` | `100 ms` | Duration | The initial delay before retrying a write that failed with a retriable error. The delay doubles with the batch retry count and uses 20% jitter. Set to `0` to disable this backoff. | +| `client.writer.retry-backoff-max` | `1 s` | Duration | The maximum retry delay. Must be between the initial backoff and `2147483647ms`. | + +The retry backoff applies per physical table bucket to both Log and KV writes. It covers transient +failures such as `NOT_ENOUGH_REPLICAS_EXCEPTION`, preventing a re-enqueued batch from being resent +on every sender cycle while the TabletServer is still rejecting writes. Other buckets remain +eligible for normal scheduling. + +The backoff is installed before the failed batch is re-enqueued. It uses the batch retry count for +exponential growth, preserves the longest active deadline when several failures race for the same +bucket, and remains subject to `client.writer.retries`. Setting `client.writer.retry-backoff` to `0` +restores immediate retry behavior. + ## Client disk write protection backoff | Key | Default | Type | Description | @@ -239,12 +256,13 @@ and `2147483647ms`. The backoff doubles with the existing batch retry count, wit and a cap at the maximum. Setting the two values equal selects a fixed delay without jitter. Actual delays are always at least `1ms`. -Disk backoff is independent of `client.writer.kv-backpressure.max-throttle`; when both apply, -the writer waits for the longer remaining duration. A successful in-flight write or a KV -pressure value of zero does not clear an active disk wait. Writes retry automatically after -the wait, subject to `client.writer.retries`. Flush and graceful close respect the wait; -the existing close timeout still applies. After disk protection is lifted, the next retry -can take up to the remaining backoff window, plus normal scheduling and RPC time. +Disk backoff is independent of `client.writer.retry-backoff` and +`client.writer.kv-backpressure.max-throttle`; when several gates apply, the writer waits for the +longest remaining duration. A successful in-flight write or a KV pressure value of zero does not +clear an active disk wait. Writes retry automatically after the wait, subject to +`client.writer.retries`. Flush and graceful close respect the wait; the existing close timeout +still applies. After disk protection is lifted, the next retry can take up to the remaining +backoff window, plus normal scheduling and RPC time. This behavior requires upgrading the Java client or the connector that bundles it. It uses the existing server error code and requires no server protocol upgrade. From 75d14779c9a2c987d1d8e71dd55d475f49ebe497 Mon Sep 17 00:00:00 2001 From: fhan Date: Mon, 5 Oct 2026 20:28:47 +0800 Subject: [PATCH 5/6] [client] Deduplicate metadata invalidation before write retries --- .../org/apache/fluss/client/write/Sender.java | 68 ++++--- .../apache/fluss/client/write/SenderTest.java | 191 +++++++++++++++--- 2 files changed, 205 insertions(+), 54 deletions(-) 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 34957b0ab0..1271fdf0c4 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 @@ -49,7 +49,6 @@ import javax.annotation.concurrent.GuardedBy; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -600,6 +599,7 @@ private void handleProduceLogResponse( long tableId, Map writeBatchesByKey) { Set invalidMetadataTablesSet = new HashSet<>(); + List batchesToRetry = new ArrayList<>(); for (PbProduceLogRespForBucket logRespForBucket : response.getBucketsRespsList()) { TableBucket tb = new TableBucket( @@ -616,15 +616,16 @@ private void handleProduceLogResponse( ? logRespForBucket.getOriginalPartitionName() : null)); if (logRespForBucket.hasErrorCode()) { - Set invalidMetadataTables = - handleWriteBatchException( - writeBatch, ApiError.fromErrorMessage(logRespForBucket)); - invalidMetadataTablesSet.addAll(invalidMetadataTables); + handleWriteBatchException( + writeBatch, + ApiError.fromErrorMessage(logRespForBucket), + invalidMetadataTablesSet, + batchesToRetry); } else { completeBatch(writeBatch); } } - metadataUpdater.invalidPhysicalTableBucketMeta(invalidMetadataTablesSet); + invalidateMetadataAndRetry(invalidMetadataTablesSet, batchesToRetry); } private void handlePutKvResponse( @@ -632,6 +633,7 @@ private void handlePutKvResponse( long tableId, Map writeBatchesByKey) { Set invalidMetadataTablesSet = new HashSet<>(); + List batchesToRetry = new ArrayList<>(); boolean eligibilityAdvanced = false; for (PbPutKvRespForBucket respForBucket : putKvResponse.getBucketsRespsList()) { TableBucket tb = @@ -656,10 +658,11 @@ private void handlePutKvResponse( continue; } if (respForBucket.hasErrorCode()) { - Set invalidMetadataTables = - handleWriteBatchException( - writeBatch, ApiError.fromErrorMessage(respForBucket)); - invalidMetadataTablesSet.addAll(invalidMetadataTables); + handleWriteBatchException( + writeBatch, + ApiError.fromErrorMessage(respForBucket), + invalidMetadataTablesSet, + batchesToRetry); } else { // Get log end offset (LEO) from response for KV batches long logEndOffset = @@ -667,7 +670,7 @@ private void handlePutKvResponse( completeBatch(writeBatch, tb, logEndOffset); } } - metadataUpdater.invalidPhysicalTableBucketMeta(invalidMetadataTablesSet); + invalidateMetadataAndRetry(invalidMetadataTablesSet, batchesToRetry); // A fresher latest-wins KV pressure signal may shorten the effective gate. Wake only after // the whole response is applied so the sender observes completed batches and invalidated // metadata together with the new deadline. @@ -688,12 +691,23 @@ private void handleWriteRequestException(Throwable t, List writ // if batch failed because of retrievable exception, we need to retry send all those // batches. Set invalidMetadataTablesSet = new HashSet<>(); + List batchesToRetry = new ArrayList<>(); for (ReadyWriteBatch batch : writeBatches) { - Set invalidMetadataTables = handleWriteBatchException(batch, error); - invalidMetadataTablesSet.addAll(invalidMetadataTables); + handleWriteBatchException(batch, error, invalidMetadataTablesSet, batchesToRetry); } - metadataUpdater.invalidPhysicalTableBucketMeta(invalidMetadataTablesSet); + invalidateMetadataAndRetry(invalidMetadataTablesSet, batchesToRetry); + } + + private void invalidateMetadataAndRetry( + Set invalidMetadataTables, List batchesToRetry) { + // Rebuild the cluster snapshot once per response/request, even when several failed + // batches share a target. Invalidate all targets before publishing any retry, since + // re-enqueueing wakes the sender and must not expose stale metadata to it. + metadataUpdater.invalidPhysicalTableBucketMeta(invalidMetadataTables); + for (ReadyWriteBatch batch : batchesToRetry) { + reEnqueueBatch(batch); + } } /** Stops appends and sending after a fatal write or partition-creation failure. */ @@ -710,10 +724,12 @@ void recordFatalError(Throwable t) { wakeup(); } - /** Handle the exception and return a set of tables for which the metadata is invalid. */ - private Set handleWriteBatchException( - ReadyWriteBatch readyWriteBatch, ApiError error) { - Set invalidMetadataTables = new HashSet<>(); + /** Handles an exception, collecting metadata invalidations and deferring retry publication. */ + private void handleWriteBatchException( + ReadyWriteBatch readyWriteBatch, + ApiError error, + Set invalidMetadataTables, + List batchesToRetry) { WriteBatch writeBatch = readyWriteBatch.writeBatch(); // Historical queues use the original path as their accumulator key, so capture the actual // RPC target before any retry handling. @@ -731,7 +747,7 @@ private Set handleWriteBatchException( // the same ordered stream. Stop accepting/draining writes and fail the writer. recordFatalError(error.exception()); invalidMetadataTables.add(writeBatch.physicalTablePath()); - return invalidMetadataTables; + return; } if (error.error() == Errors.DUPLICATE_SEQUENCE_EXCEPTION) { // If we have received a duplicate batch sequence error, it means that the batch @@ -768,16 +784,13 @@ private Set handleWriteBatchException( readyWriteBatch.tableBucket(), error.exception()); } - // Re-enqueuing publishes the retry to the sender and wakes it up. Invalidate the - // actual RPC target first so the retry cannot race ahead using stale metadata. A - // historical batch remains keyed by its original partition path in the - // accumulator, while its RPC is sent to the internal historical partition. - metadataUpdater.invalidPhysicalTableBucketMeta( - Collections.singleton(writeTargetPath)); + // Historical batches keep their original accumulator path, but metadata + // invalidation must target the physical partition used by the RPC. + invalidMetadataTables.add(writeTargetPath); } if (!idempotenceManager.idempotenceEnabled()) { prepareWriteRetry(readyWriteBatch, error); - reEnqueueBatch(readyWriteBatch); + batchesToRetry.add(readyWriteBatch); } else if (idempotenceManager.hasWriterId(writeBatch.writerId())) { // If idempotence is enabled only retry the request if the current writer id is // the same as the writer id of the batch. @@ -786,7 +799,7 @@ private Set handleWriteBatchException( readyWriteBatch.tableBucket(), writeBatch.batchSequence()); prepareWriteRetry(readyWriteBatch, error); - reEnqueueBatch(readyWriteBatch); + batchesToRetry.add(readyWriteBatch); } else { Exception exception = Errors.UNKNOWN_WRITER_ID_EXCEPTION.exception( @@ -806,7 +819,6 @@ private Set handleWriteBatchException( // sequence was accepted or not, and thus it is not safe to reassign the sequence. failBatch(readyWriteBatch, error.exception(), writeBatch.attempts() < this.retries); } - return invalidMetadataTables; } private void prepareWriteRetry(ReadyWriteBatch batch, ApiError error) { 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 ad6bde5221..adc099bf9c 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 @@ -244,42 +244,181 @@ void testReroutesWriteAfterExplicitMissingPartitionResponse(boolean historicalCo } @ParameterizedTest - @ValueSource(booleans = {false, true}) - void testInvalidatesMetadataBeforePublishingRetry(boolean kv) throws Exception { + @CsvSource({"false,false", "false,true", "true,false", "true,true"}) + void testInvalidatesMetadataOnceBeforePublishingRetries(boolean kv, boolean rpcFailure) + throws Exception { sender.destroyResources(); - Map tableInfos = new HashMap<>(); - tableInfos.put(DATA1_TABLE_PATH, DATA1_TABLE_INFO); - tableInfos.put(DATA1_TABLE_PATH_PK, DATA1_TABLE_INFO_PK); - AtomicReference retryQueuedAtInvalidation = new AtomicReference<>(); + TableInfo tableInfo = kv ? DATA1_TABLE_INFO_PK : DATA1_TABLE_INFO; + PhysicalTablePath path = PhysicalTablePath.of(tableInfo.getTablePath()); + List> invalidations = new ArrayList<>(); + List retriesQueuedAtInvalidation = new ArrayList<>(); metadataUpdater = - new TestingMetadataUpdater(tableInfos) { - @Override - public void invalidPhysicalTableBucketMeta( - Set physicalTablesToInvalid) { - if (!physicalTablesToInvalid.isEmpty()) { - retryQueuedAtInvalidation.set(accumulator.hasUnDrained()); - } - super.invalidPhysicalTableBucketMeta(physicalTablesToInvalid); - } - }; + metadataUpdaterRecordingInvalidations( + tableInfo, invalidations, retriesQueuedAtInvalidation); + TableBucket firstBucket = new TableBucket(tableInfo.getTableId(), 0); + TableBucket secondBucket = new TableBucket(tableInfo.getTableId(), 1); + Cluster cluster = + new Cluster( + Collections.singletonMap( + TestingMetadataUpdater.NODE1.id(), TestingMetadataUpdater.NODE1), + TestingMetadataUpdater.COORDINATOR, + Collections.singletonMap( + path, + Arrays.asList( + new BucketLocation(path, firstBucket, 1, new int[] {1}), + new BucketLocation(path, secondBucket, 1, new int[] {1}))), + Collections.singletonMap(tableInfo.getTablePath(), tableInfo.getTableId()), + Collections.emptyMap(), + Collections.singletonMap( + TableOrPartition.ofTable(tableInfo.getTableId()), + tableInfo.getNumBuckets())); + metadataUpdater.updateCluster(cluster); sender = setupWithIdempotenceState(); - TableBucket bucket = kv ? new TableBucket(DATA1_TABLE_ID_PK, 0) : tb1; - CompletableFuture result = appendDiskTestRecord(kv, bucket, 1); + CompletableFuture firstResult = appendDiskTestRecord(kv, firstBucket, 1); + CompletableFuture secondResult = appendDiskTestRecord(kv, secondBucket, 2); sender.runOnce(); - finishRequest( - tb1, - 0, - kv - ? createPutKvResponse(bucket, Errors.NOT_LEADER_OR_FOLLOWER) - : createProduceLogResponse(bucket, Errors.NOT_LEADER_OR_FOLLOWER)); + TestTabletServerGateway gateway = node1Gateway(); + assertThat(gateway.pendingRequestSize()).isOne(); + ApiMessage request = gateway.getRequest(0); + assertThat( + kv + ? ((PutKvRequest) request).getBucketsReqsCount() + : ((ProduceLogRequest) request).getBucketsReqsCount()) + .isEqualTo(2); + if (rpcFailure) { + gateway.failRequest(0, Errors.NOT_LEADER_OR_FOLLOWER.exception()); + } else { + gateway.response( + 0, + kv + ? makePutKvResponse( + Arrays.asList( + new PutKvResultForBucket( + firstBucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError()), + new PutKvResultForBucket( + secondBucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError()))) + : makeProduceLogResponse( + Arrays.asList( + new ProduceLogResultForBucket( + firstBucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError()), + new ProduceLogResultForBucket( + secondBucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError())))); + } - assertThat(retryQueuedAtInvalidation.get()).isFalse(); - assertThat(result).isNotDone(); + assertThat(invalidations).containsExactly(Collections.singleton(path)); + assertThat(retriesQueuedAtInvalidation).containsExactly(false); + assertThat(firstResult).isNotDone(); + assertThat(secondResult).isNotDone(); + assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isEqualTo(2); + assertThat(metadataUpdater.getCluster().getBucketLocation(firstBucket)).isEmpty(); + assertThat(metadataUpdater.getCluster().getBucketLocation(secondBucket)).isEmpty(); assertThat(accumulator.hasUnDrained()).isTrue(); accumulator.abortAllBatches(new RuntimeException("test cleanup")); } + @ParameterizedTest + @CsvSource({"false,false", "false,true", "true,false", "true,true"}) + void testInvalidatesHistoricalTargetOnceBeforePublishingRetries(boolean kv, boolean rpcFailure) + throws Exception { + sender.destroyResources(); + TableInfo info = createHistoricalTableInfo(AutoPartitionTimeUnit.DAY, 7, kv); + PhysicalTablePath firstOriginal = PhysicalTablePath.of(info.getTablePath(), "20000101"); + PhysicalTablePath secondOriginal = PhysicalTablePath.of(info.getTablePath(), "20000102"); + PhysicalTablePath historical = + PhysicalTablePath.of(info.getTablePath(), HISTORICAL_PARTITION_VALUE); + TableBucket bucket = new TableBucket(info.getTableId(), 22L, 0); + List> invalidations = new ArrayList<>(); + List retriesQueuedAtInvalidation = new ArrayList<>(); + metadataUpdater = + metadataUpdaterRecordingInvalidations( + info, invalidations, retriesQueuedAtInvalidation); + metadataUpdater.updateCluster( + partitionedCluster(info, Collections.singletonMap(historical, bucket))); + sender = setupWithIdempotenceState(); + accumulator.checkAndCacheHistoricalPartitionEnabled(info); + List> results = new ArrayList<>(); + for (PhysicalTablePath original : Arrays.asList(firstOriginal, secondOriginal)) { + accumulator.routeWritesTo(info, original, historical, bucket.getPartitionId()); + if (kv) { + results.add(appendKvRecord(info, original, 1, metadataUpdater.getCluster())); + } else { + CompletableFuture result = new CompletableFuture<>(); + results.add(result); + bucketAssigner.setBucketId(0); + accumulator.append( + WriteRecord.forArrowAppend( + info, original, row(1, original.getPartitionName()), null), + (tb, offset, error) -> result.complete(error), + metadataUpdater.getCluster()); + } + } + sender.runOnce(); + TestTabletServerGateway gateway = node1Gateway(); + assertThat(gateway.pendingRequestSize()).isOne(); + ApiMessage request = gateway.getRequest(0); + assertThat( + kv + ? ((PutKvRequest) request).getBucketsReqsCount() + : ((ProduceLogRequest) request).getBucketsReqsCount()) + .isEqualTo(2); + if (rpcFailure) { + gateway.failRequest(0, Errors.NOT_LEADER_OR_FOLLOWER.exception()); + } else { + gateway.response( + 0, + kv + ? makePutKvResponse( + Arrays.asList( + PutKvResultForBucket.historicalFailure( + bucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError(), + firstOriginal.getPartitionName()), + PutKvResultForBucket.historicalFailure( + bucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError(), + secondOriginal.getPartitionName()))) + : makeProduceLogResponse( + Arrays.asList( + ProduceLogResultForBucket.historicalFailure( + bucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError(), + firstOriginal.getPartitionName()), + ProduceLogResultForBucket.historicalFailure( + bucket, + Errors.NOT_LEADER_OR_FOLLOWER.toApiError(), + secondOriginal.getPartitionName())))); + } + assertThat(invalidations).containsExactly(Collections.singleton(historical)); + assertThat(retriesQueuedAtInvalidation).containsExactly(false); + results.forEach(result -> assertThat(result).isNotDone()); + assertThat(writerMetricGroup.recordsRetryTotal().getCount()).isEqualTo(2); + assertThat(metadataUpdater.getCluster().getBucketLocation(bucket)).isEmpty(); + assertThat(accumulator.hasUnDrained()).isTrue(); + accumulator.abortAllBatches(new RuntimeException("test cleanup")); + } + + private TestingMetadataUpdater metadataUpdaterRecordingInvalidations( + TableInfo info, + List> invalidations, + List retriesQueuedAtInvalidation) { + return new TestingMetadataUpdater(Collections.singletonMap(info.getTablePath(), info)) { + @Override + public void invalidPhysicalTableBucketMeta( + Set physicalTablesToInvalid) { + if (!physicalTablesToInvalid.isEmpty()) { + invalidations.add(new HashSet<>(physicalTablesToInvalid)); + retriesQueuedAtInvalidation.add(accumulator.hasUnDrained()); + } + super.invalidPhysicalTableBucketMeta(physicalTablesToInvalid); + } + }; + } + @ParameterizedTest @CsvSource({"2, 4", "4, 2"}) void testServerValidatesResolvedBucketCount(int initialCount, int updatedCount) From 6798b47668513af3b29a9709180a8f56f828233b Mon Sep 17 00:00:00 2001 From: fhan Date: Wed, 7 Oct 2026 20:33:56 +0800 Subject: [PATCH 6/6] [flink] Release pending tiering rounds when closing enumerator --- .../enumerator/TieringSourceEnumerator.java | 166 +++++++++++------- .../TieringSourceEnumeratorTest.java | 90 ++++++++++ 2 files changed, 189 insertions(+), 67 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java index 2f30a08ee5..dec852a376 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java @@ -119,6 +119,10 @@ public class TieringSourceEnumerator private final Map finishedTables; private final Set tieringReachMaxDurationsTables; + // Serialize heartbeat acknowledgements and table-state transitions with shutdown. An + // in-flight heartbeat may acquire a new table that the final heartbeat must also release. + private final Object lifecycleLock = new Object(); + // lazily instantiated private RpcClient rpcClient; private CoordinatorGateway coordinatorGateway; @@ -301,6 +305,15 @@ private int max(Set integers) { @Override public void handleSourceEvent(int subtaskId, SourceEvent sourceEvent) { + synchronized (lifecycleLock) { + if (closed) { + return; + } + handleTieringEvent(sourceEvent); + } + } + + private void handleTieringEvent(SourceEvent sourceEvent) { if (sourceEvent instanceof FinishedTieringEvent) { FinishedTieringEvent finishedTieringEvent = (FinishedTieringEvent) sourceEvent; long finishedTableId = finishedTieringEvent.getTableId(); @@ -346,18 +359,23 @@ public void handleSourceEvent(int subtaskId, SourceEvent sourceEvent) { } private void handleSourceReaderFailOver() { - LOG.info( - "Handling source reader fail over, mark current tiering table epoch {} as failed.", - tieringTableEpochs); - // we need to make all as failed - failedTableEpochs.putAll(new HashMap<>(tieringTableEpochs)); - tieringTableEpochs.clear(); - tieringReachMaxDurationsTables.clear(); - // also clean all pending splits since we mark all as failed - pendingSplits.clear(); - if (!failedTableEpochs.isEmpty()) { - // call one round of heartbeat to notify table has been finished or failed - requestTableAndAssign(0); + synchronized (lifecycleLock) { + if (closed) { + return; + } + LOG.info( + "Handling source reader fail over, mark current tiering table epoch {} as failed.", + tieringTableEpochs); + // we need to make all as failed + failedTableEpochs.putAll(new HashMap<>(tieringTableEpochs)); + tieringTableEpochs.clear(); + tieringReachMaxDurationsTables.clear(); + // also clean all pending splits since we mark all as failed + pendingSplits.clear(); + if (!failedTableEpochs.isEmpty()) { + // call one round of heartbeat to notify table has been finished or failed + requestTableAndAssign(0); + } } } @@ -392,18 +410,23 @@ protected void handleTableTieringReachMaxDuration( @VisibleForTesting void generateAndAssignSplits( @Nullable Tuple3 tieringTable, Throwable throwable) { - if (throwable != null) { - ExceptionUtils.rethrow(throwable); - } - if (tieringTable != null) { - generateTieringSplits(tieringTable); + synchronized (lifecycleLock) { + if (closed) { + return; + } + if (throwable != null) { + ExceptionUtils.rethrow(throwable); + } + if (tieringTable != null) { + generateTieringSplits(tieringTable); + } + assignSplits(); } - assignSplits(); } private void assignSplits() { - // we don't assign splits during failover - if (isFailOvering) { + // Do not assign splits during failover or after shutdown. + if (closed || isFailOvering) { return; } if (!readersAwaitingSplit.isEmpty()) { @@ -424,6 +447,12 @@ private void assignSplits() { } private @Nullable Tuple3 requestTieringTableSplitsViaHeartBeat() { + synchronized (lifecycleLock) { + return requestTieringTableSplits(); + } + } + + private @Nullable Tuple3 requestTieringTableSplits() { if (closed) { return null; } @@ -581,58 +610,61 @@ public TieringSourceEnumeratorState snapshotState(long checkpointId) throws Exce @Override public void close() throws IOException { - closed = true; - timerService.shutdownNow(); - if (rpcClient != null) { - failedTableEpochs.putAll(tieringTableEpochs); - tieringTableEpochs.clear(); - if (!failedTableEpochs.isEmpty()) { - reportFailedTable(basicHeartBeat(), failedTableEpochs); + synchronized (lifecycleLock) { + if (closed) { + return; + } + closed = true; + timerService.shutdownNow(); + if (rpcClient != null) { + failedTableEpochs.putAll(tieringTableEpochs); + tieringTableEpochs.clear(); + try { + // Empty split rounds and finished events may still be waiting for the next + // heartbeat. Report them as finished (including force-finish and stats), while + // releasing unfinished rounds as failed. Never request another table on close. + if (!finishedTables.isEmpty() || !failedTableEpochs.isEmpty()) { + waitHeartbeatResponse( + coordinatorGateway.lakeTieringHeartbeat( + tieringTableHeartBeat( + basicHeartBeat(), + Collections.emptyMap(), + finishedTables, + failedTableEpochs, + flussCoordinatorEpoch))); + finishedTables.clear(); + failedTableEpochs.clear(); + } + } catch (Exception e) { + // The coordinator timeout remains the fallback if it cannot be notified. + // Reporting failure must not prevent the clients from being closed. + LOG.warn("Failed to report final tiering table states during close.", e); + } + try { + LOG.info( + "Closing Tiering Source Enumerator of at {}.", + System.currentTimeMillis()); + rpcClient.close(); + } catch (Exception e) { + LOG.error("Failed to close Tiering Source enumerator.", e); + } } try { - LOG.info("Closing Tiering Source Enumerator of at {}.", System.currentTimeMillis()); - rpcClient.close(); + if (flussAdmin != null) { + LOG.info("Closing Fluss Admin client..."); + flussAdmin.close(); + } } catch (Exception e) { - LOG.error("Failed to close Tiering Source enumerator.", e); + LOG.error("Failed to close Fluss Admin client.", e); } - } - try { - if (flussAdmin != null) { - LOG.info("Closing Fluss Admin client..."); - flussAdmin.close(); - } - } catch (Exception e) { - LOG.error("Failed to close Fluss Admin client.", e); - } - try { - if (connection != null) { - LOG.info("Closing Fluss connection..."); - connection.close(); + try { + if (connection != null) { + LOG.info("Closing Fluss connection..."); + connection.close(); + } + } catch (Exception e) { + LOG.error("Failed to close Fluss connection.", e); } - } catch (Exception e) { - LOG.error("Failed to close Fluss connection.", e); - } - } - - /** - * Report failed table to Fluss coordinator via HeartBeat, this method should be called when - * {@link TieringSourceEnumerator} is closed or receives failed table from downstream lake - * committer. - */ - private void reportFailedTable( - LakeTieringHeartbeatRequest heartbeatRequest, Map failedTableEpochs) - throws FlinkRuntimeException { - try { - waitHeartbeatResponse( - coordinatorGateway.lakeTieringHeartbeat( - failedTableHeartBeat( - heartbeatRequest, failedTableEpochs, flussCoordinatorEpoch))); - LOG.info("Report failed table to Fluss Coordinator success"); - - } catch (Exception e) { - LOG.error("Errors happens when report failed table to Fluss cluster.", e); - throw new FlinkRuntimeException( - "Errors happens when report failed table to Fluss cluster.", e); } } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java index c91a960e95..3a8c2c8570 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java @@ -42,10 +42,13 @@ import org.apache.flink.api.connector.source.SourceEvent; import org.apache.flink.api.connector.source.SplitsAssignment; import org.apache.flink.api.connector.source.mocks.MockSplitEnumeratorContext; +import org.apache.flink.api.java.tuple.Tuple3; import org.apache.flink.util.FlinkRuntimeException; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import javax.annotation.Nullable; @@ -56,6 +59,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -82,6 +86,92 @@ protected void beforeEach() throws Exception { super.beforeEach(); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testCloseReportsEmptyRoundBeforeNextHeartbeat(boolean primaryKey) throws Throwable { + TablePath tablePath = TablePath.of(DEFAULT_DB, "close-empty-round-" + primaryKey); + TableDescriptor descriptor = + primaryKey ? DEFAULT_PK_TABLE_DESCRIPTOR : DEFAULT_LOG_TABLE_DESCRIPTOR; + long tableId = createTable(tablePath, descriptor); + AtomicInteger generatedRounds = new AtomicInteger(); + TestingLakeTieringFactory factory = + new TestingLakeTieringFactory() { + @Override + public void validateTable(TableInfo tableInfo) { + generatedRounds.incrementAndGet(); + } + }; + + try (FlussMockSplitEnumeratorContext context = + new FlussMockSplitEnumeratorContext<>(1); + TieringSourceEnumerator enumerator = + createTieringSourceEnumerator(flussConf, context, factory)) { + enumerator.start(); + registerSingleReaderAndHandleSplitRequests(context, enumerator, 0, 0); + assertThat(generatedRounds.get()).isOne(); + assertThat(context.getSplitsAssignmentSequence()).isEmpty(); + + // Do not run another heartbeat: the empty round has only been finished locally. + enumerator.close(); + + // Queued async work must not acquire a new round or access closed clients. + context.runPeriodicCallable(0); + enumerator.generateAndAssignSplits(Tuple3.of(tableId, 1L, tablePath), null); + assertThat(generatedRounds.get()).isOne(); + assertThat(context.getSplitsAssignmentSequence()).isEmpty(); + } + + // The next service must tier newly written data without waiting for the two-minute + // coordinator timeout that would otherwise reclaim the abandoned empty round. + if (primaryKey) { + upsertRow(tablePath, descriptor, 0, 10); + } else { + appendRow(tablePath, descriptor, 0, 10); + } + assertNextEnumeratorCanAcquireTable(tablePath); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testCloseReleasesActiveOrFinishedRound(boolean finished) throws Exception { + TablePath tablePath = + TablePath.of(DEFAULT_DB, "close-active-or-finished-round-" + finished); + long tableId = createTable(tablePath, DEFAULT_LOG_TABLE_DESCRIPTOR); + appendRow(tablePath, DEFAULT_LOG_TABLE_DESCRIPTOR, 0, 10); + try (FlussMockSplitEnumeratorContext context = + new FlussMockSplitEnumeratorContext<>(1); + TieringSourceEnumerator enumerator = + createTieringSourceEnumerator(flussConf, context)) { + enumerator.start(); + registerSingleReaderAndHandleSplitRequests(context, enumerator, 0, 0); + assertThat(context.getSplitsAssignmentSequence()).isNotEmpty(); + if (finished) { + // Queue a completion report, but close before its async heartbeat executes. + enumerator.handleSourceEvent(0, new FinishedTieringEvent(tableId)); + } + } + assertNextEnumeratorCanAcquireTable(tablePath); + } + + private void assertNextEnumeratorCanAcquireTable(TablePath tablePath) throws Exception { + try (FlussMockSplitEnumeratorContext context = + new FlussMockSplitEnumeratorContext<>(1); + TieringSourceEnumerator enumerator = + createTieringSourceEnumerator(flussConf, context)) { + enumerator.start(); + context.registerSourceReader(0, 0, "localhost-0"); + enumerator.addReader(0); + retry( + Duration.ofSeconds(15), + () -> { + enumerator.handleSplitRequest(0, "localhost-0"); + assertThat(context.getSplitsAssignmentSequence()).isNotEmpty(); + }); + assertThat(context.getSplitsAssignmentSequence().get(0).assignment().get(0)) + .allSatisfy(split -> assertThat(split.getTablePath()).isEqualTo(tablePath)); + } + } + @Test void testPrimaryKeyTableWithNoSnapshotSplits() throws Throwable { TablePath tablePath = DEFAULT_TABLE_PATH;