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..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 @@ -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,19 +137,11 @@ 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; + // 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; - // 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(); - - // 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 @@ -159,6 +155,8 @@ public final class RecordAccumulator { idempotenceManager, writerMetricGroup, clock, + createRetriableWriteBackoff(conf), + createDiskWriteBackoff(conf), BucketAssignerFactory.defaultFactory(conf)); } @@ -168,6 +166,60 @@ public final class RecordAccumulator { IdempotenceManager idempotenceManager, WriterMetricGroup writerMetricGroup, Clock clock, + ExponentialBackoff diskWriteBackoff) { + this( + conf, + idempotenceManager, + writerMetricGroup, + clock, + createRetriableWriteBackoff(conf), + diskWriteBackoff, + BucketAssignerFactory.defaultFactory(conf)); + } + + @VisibleForTesting + RecordAccumulator( + Configuration conf, + IdempotenceManager idempotenceManager, + WriterMetricGroup writerMetricGroup, + Clock clock, + BucketAssignerFactory bucketAssignerFactory) { + this( + conf, + idempotenceManager, + writerMetricGroup, + clock, + createRetriableWriteBackoff(conf), + createDiskWriteBackoff(conf), + bucketAssignerFactory); + } + + @VisibleForTesting + RecordAccumulator( + Configuration conf, + IdempotenceManager idempotenceManager, + WriterMetricGroup writerMetricGroup, + 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; @@ -193,11 +245,45 @@ 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(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); + 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 +392,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 +407,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 +722,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 +856,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 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 +1349,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 +1390,56 @@ 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); } /** - * 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. + * Installs general retry backoff before the batch retry count is increased by re-enqueueing. */ - 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); + 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()); + } + + /** Returns the remaining disk backoff without changing its deadline. */ + long diskWriteBackoffRemainingMs(TableBucket tableBucket) { + return throttle.diskRemainingMs(tableBucket); + } + + /** Reclaims expired backoff entries even when their queues no longer contain any batches. */ + void maybeEvictExpiredWriteBackoffs() { + throttle.maybeEvictExpiredBackoffs(); + } + + @VisibleForTesting + int retriableWriteBackoffCount() { + return throttle.retryBackoffCount(); + } + + @VisibleForTesting + int diskWriteBackoffCount() { + return throttle.diskBackoffCount(); + } + + boolean updateThrottle(TableBucket tableBucket, float pressure) { + return 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 +1654,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.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 9f8b1e8ba2..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 @@ -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.maybeEvictExpiredWriteBackoffs(); // 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 @@ -597,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( @@ -613,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( @@ -629,6 +633,8 @@ 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 = new TableBucket( @@ -638,7 +644,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 = @@ -652,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 = @@ -663,7 +670,13 @@ 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. + if (eligibilityAdvanced) { + wakeup(); + } } private void handleWriteRequestException(Throwable t, List writeBatches) { @@ -678,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. */ @@ -700,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. @@ -721,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 @@ -746,14 +772,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. - LOG.warn( - "Get error write response on table bucket {}, retrying ({} attempts left). Error: {}", - readyWriteBatch.tableBucket(), - retries - writeBatch.attempts(), - error.formatErrMsg()); - + 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()); + } + // 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()) { - reEnqueueBatch(readyWriteBatch); + prepareWriteRetry(readyWriteBatch, error); + 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. @@ -761,7 +798,8 @@ private Set handleWriteBatchException( "Retrying batch to table-bucket {}, Batch sequence : {}", readyWriteBatch.tableBucket(), writeBatch.batchSequence()); - reEnqueueBatch(readyWriteBatch); + prepareWriteRetry(readyWriteBatch, error); + batchesToRetry.add(readyWriteBatch); } else { Exception exception = Errors.UNKNOWN_WRITER_ID_EXCEPTION.exception( @@ -771,24 +809,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: {}", @@ -799,7 +819,29 @@ 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) { + long retryBackoffMs = accumulator.backoffAfterRetriableWrite(batch); + if (error.error() == Errors.DISK_WRITE_LOCKED) { + long diskBackoffMs = accumulator.backoffAfterDiskWriteLocked(batch); + LOG.warn( + "Get error write response on table bucket {}, disk backoff {} ms and " + + "retry backoff {} ms ({} attempts left). Error: {}", + batch.tableBucket(), + diskBackoffMs, + retryBackoffMs, + retries - batch.writeBatch().attempts(), + error.formatErrMsg()); + } else { + LOG.warn( + "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 new file mode 100644 index 0000000000..3183181ae6 --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/write/WriteThrottleController.java @@ -0,0 +1,332 @@ +/* + * 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, + * retriable-write backoff, and disk-write backoff. All gates are per-{@link TableBucket}. + * + *

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. + *
+ * + *

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; + + // 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 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. + 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 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.lastBackoffSweepNanos = 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. 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), + Math.max(retryBackoffRemainingMs(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. + * @return whether this update moved the effective eligibility time earlier, after accounting + * for the KV, retriable-write, and disk gates + */ + 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. + newKvRemainingMs = maxThrottleMs; + } else if (pressure > 0f) { + newKvRemainingMs = (long) (maxThrottleMs * pressure * pressure); + } + + 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); + } + + long previousKvRemainingMs = + previousExpiryMs == null ? 0 : Math.max(0, previousExpiryMs - nowMs); + long otherRemainingMs = + Math.max(retryBackoffRemainingMs(tableBucket), diskRemainingMs(tableBucket)); + long previousEffectiveRemainingMs = Math.max(previousKvRemainingMs, otherRemainingMs); + long newEffectiveRemainingMs = Math.max(newKvRemainingMs, otherRemainingMs); + return newEffectiveRemainingMs < previousEffectiveRemainingMs; + } + + 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); + } + + // ------------------------------------------------------------------------ + // 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). + // ------------------------------------------------------------------------ + + /** + * 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) { + 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(delayMs); + Long deadline = + deadlines.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()) { + deadlines.remove(tableBucket, deadline); + return 0; + } + return nanosToCeilMillis(deadline - now); + } + + private long remainingBackoffMs( + ConcurrentMap deadlines, TableBucket tableBucket) { + Long deadline = deadlines.get(tableBucket); + if (deadline == null) { + return 0; + } + long remainingNanos = deadline - clock.nanoseconds(); + if (remainingNanos > 0) { + return nanosToCeilMillis(remainingNanos); + } + deadlines.remove(tableBucket, deadline); + return 0; + } + + /** 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 - lastBackoffSweepNanos < TimeUnit.SECONDS.toNanos(1)) { + return; + } + lastBackoffSweepNanos = now; + evictExpiredBackoffs(retryBackoffDeadlineNanos, now); + evictExpiredBackoffs(diskBackoffDeadlineNanos, now); + } + + private static void evictExpiredBackoffs(ConcurrentMap deadlines, long now) { + deadlines.forEach( + (bucket, deadline) -> { + if (deadline - now <= 0) { + deadlines.remove(bucket, deadline); + } + }); + } + + /** 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(); + } + + 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..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 @@ -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,293 @@ 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 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)); + 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 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 two-second + // general retry deadline. + assertThat(accum.updateThrottle(tb1, 0.1f)).isTrue(); + 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(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(); + } + + @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::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::maybeEvictExpiredWriteBackoffs), + 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.maybeEvictExpiredWriteBackoffs(); + 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"); + } + } + + @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)); + 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 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( + 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..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 @@ -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. @@ -236,6 +243,182 @@ void testReroutesWriteAfterExplicitMissingPartitionResponse(boolean historicalCo assertThat(future.get()).isNull(); } + @ParameterizedTest + @CsvSource({"false,false", "false,true", "true,false", "true,true"}) + void testInvalidatesMetadataOnceBeforePublishingRetries(boolean kv, boolean rpcFailure) + throws Exception { + sender.destroyResources(); + 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 = + 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(); + CompletableFuture firstResult = appendDiskTestRecord(kv, firstBucket, 1); + CompletableFuture secondResult = appendDiskTestRecord(kv, secondBucket, 2); + 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( + 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(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) @@ -334,8 +517,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 +581,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( @@ -1574,6 +1796,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); @@ -1618,16 +1892,22 @@ 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)); 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 +1942,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 +1962,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 +1983,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 +2003,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 +2097,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 +2386,301 @@ 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 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 { + 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)); } @@ -2081,18 +2691,29 @@ 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( 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..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,47 @@ 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() + .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/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; 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 --> +