From 60ef3800537b9d8318408dfa2a8f28b6a1592867 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=96=B9=E6=99=93=E5=85=B5?= Date: Mon, 3 Aug 2026 17:33:32 +0800 Subject: [PATCH 1/3] [server] Roll expired active segments for remote retention - Roll non-empty expired active segments after TTL once the high watermark reaches the log end offset. - Preserve contiguous inactive-segment boundaries and defer deletion until remote upload completes. - Cover partitioned and non-partitioned retention behavior, including disabled TTL and remote log settings. --- .../apache/fluss/server/log/LogTablet.java | 30 +++++-- .../log/remote/TieredLocalSegmentTTLTest.java | 83 +++++++++++++++++-- 2 files changed, 102 insertions(+), 11 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index 875bef312a..f097165924 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -109,6 +109,7 @@ public final class LogTablet { private final Clock clock; private final boolean isChangeLog; private final long logTtlMs; + private final boolean remoteLogEnabled; @GuardedBy("lock") private volatile LogOffsetMetadata highWatermarkMetadata; @@ -160,6 +161,8 @@ private LogTablet( this.writerStateManager = writerStateManager; this.highWatermarkMetadata = new LogOffsetMetadata(0L); this.logTtlMs = logTtlMs; + this.remoteLogEnabled = + conf.get(ConfigOptions.REMOTE_LOG_TASK_INTERVAL_DURATION).toMillis() > 0L; this.scheduler = scheduler; // scheduler the writer expiration interval check. @@ -1048,7 +1051,7 @@ private void maybeRoll(int messageSize, LogAppendInfo appendInfo) throws Excepti */ @VisibleForTesting @SuppressWarnings("OptionalUsedAsFieldOrParameterType") - public void roll(Optional expectedNextOffset) throws Exception { + public void roll(Optional expectedNextOffset) throws IOException { synchronized (lock) { LogSegment segment = localLog.roll(expectedNextOffset); // Take a snapshot of the writer state to facilitate recovery. It is useful to have @@ -1350,7 +1353,7 @@ private List deletableRemoteSegments(long endOffset) { return deletableSegments; } - /** Returns the contiguous prefix of inactive segments that has expired. */ + /** Returns the contiguous prefix of expired segments and rolls an expired active segment. */ private List deletableExpiredSegments(long endOffset) throws IOException { if (localLog.getSegments().isEmpty()) { return Collections.emptyList(); @@ -1359,13 +1362,28 @@ private List deletableExpiredSegments(long endOffset) throws IOExcep List deletableSegments = new ArrayList<>(); List logSegments = localLog.getSegments().values(); long now = clock.milliseconds(); + boolean shouldRoll = false; - for (int i = 0; i < logSegments.size() - 1; i++) { - if (logSegments.get(i + 1).getBaseOffset() > endOffset - || !isSegmentExpired(now, logSegments.get(i), logTtlMs)) { + for (int i = 0; i < logSegments.size(); i++) { + LogSegment segment = logSegments.get(i); + boolean active = i == logSegments.size() - 1; + if ((!active && logSegments.get(i + 1).getBaseOffset() > endOffset) + || !isSegmentExpired(now, segment, logTtlMs)) { break; } - deletableSegments.add(logSegments.get(i)); + + if (active) { + shouldRoll = + remoteLogEnabled + && segment.getSizeInBytes() > 0 + && getHighWatermark() >= localLogEndOffset(); + break; + } + deletableSegments.add(segment); + } + + if (shouldRoll) { + roll(Optional.empty()); } return deletableSegments; } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java index 1ea904ef9a..3aee9e64f3 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java @@ -33,7 +33,7 @@ import static org.apache.fluss.record.TestData.DATA1_TABLE_PATH; import static org.assertj.core.api.Assertions.assertThat; -/** Tests TTL-based cleanup of inactive local segments retained by tiered storage. */ +/** Tests TTL-based cleanup of local segments retained by tiered storage. */ final class TieredLocalSegmentTTLTest extends RemoteLogTestBase { @BeforeEach @@ -71,17 +71,60 @@ void testInactiveTieredLocalSegmentRemovedAfterTtl(boolean partitionTable) throw // Run local retention without updating the remote manifest or uploading another segment. logManager.cleanupExpiredLocalLogSegments(); - // The inactive segment is expired and deleted, while the active segment is retained. + // The expired active segment is rolled, but is not deleted in the same retention pass. assertThat(remoteLog.allRemoteLogSegments()).hasSize(4); + assertThat(logTablet.getSegments()).hasSize(2); + assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); + + // An empty active segment must not be rolled again while the previous segment is waiting + // for remote upload. + logManager.cleanupExpiredLocalLogSegments(); + assertThat(logTablet.getSegments()).hasSize(2); + + remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); + assertThat(remoteLog.getRemoteLogEndOffset()).hasValue(50L); + logManager.cleanupExpiredLocalLogSegments(); + assertThat(logTablet.getSegments()).hasSize(1); + assertThat(logTablet.localLogStartOffset()).isEqualTo(50L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testExpiredActiveSegmentWaitsForHighWatermarkWithoutRemoteLogEndOffset( + boolean partitionTable) throws Exception { + TableBucket tb = + partitionTable + ? new TableBucket(DATA1_TABLE_ID, 0L, 0) + : new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(tb, partitionTable); + LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); + + addMultiSegmentsToLogTablet(logTablet, 5); + assertThat(remoteLogManager.remoteLogTablet(tb).getRemoteLogEndOffset()).isEmpty(); + + manualClock.advanceTime(Duration.ofHours(2)); + logTablet.updateHighWatermark(logTablet.localLogEndOffset() - 1L); + logManager.cleanupExpiredLocalLogSegments(); + assertThat(logTablet.getSegments()).hasSize(1); assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); + + logTablet.updateHighWatermark(logTablet.localLogEndOffset()); + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(2); + assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); } @ParameterizedTest @ValueSource(booleans = {true, false}) - void testExpiredLocalSegmentsRemovedWithoutRemoteLogEndOffset(boolean partitionTable) + void testExpiredActiveSegmentNotRolledWhenRemoteLogDisabled(boolean partitionTable) throws Exception { + conf.set(ConfigOptions.REMOTE_LOG_TASK_INTERVAL_DURATION, Duration.ZERO); TableBucket tb = partitionTable ? new TableBucket(DATA1_TABLE_ID, 0L, 0) @@ -90,8 +133,6 @@ void testExpiredLocalSegmentsRemovedWithoutRemoteLogEndOffset(boolean partitionT LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); - assertThat(remoteLogManager.remoteLogTablet(tb).getRemoteLogEndOffset()).isEmpty(); - manualClock.advanceTime(Duration.ofHours(2)); logManager.cleanupExpiredLocalLogSegments(); @@ -100,6 +141,31 @@ void testExpiredLocalSegmentsRemovedWithoutRemoteLogEndOffset(boolean partitionT assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); } + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testExpiredActiveSegmentNotRolledWhenTtlDisabled(boolean partitionTable) throws Exception { + registerTableInZkClient( + DATA1_TABLE_PATH, + DATA1_SCHEMA, + DATA1_TABLE_ID, + Collections.emptyList(), + Collections.singletonMap(ConfigOptions.TABLE_LOG_TTL.key(), "0ms")); + TableBucket tb = + partitionTable + ? new TableBucket(DATA1_TABLE_ID, 0L, 0) + : new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(tb, partitionTable); + LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); + + addMultiSegmentsToLogTablet(logTablet, 5); + manualClock.advanceTime(Duration.ofHours(2)); + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(5); + assertThat(logTablet.localLogStartOffset()).isEqualTo(0L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) void testTtlCleanupBoundedByRemoteLogEndOffset(boolean partitionTable) throws Exception { @@ -120,5 +186,12 @@ void testTtlCleanupBoundedByRemoteLogEndOffset(boolean partitionTable) throws Ex assertThat(logTablet.getSegments()).hasSize(3); assertThat(logTablet.localLogStartOffset()).isEqualTo(20L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); + + logTablet.updateRemoteLogEndOffset(40L); + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(2); + assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); } } From ccd9768e1279aea817c68cd6f02e81046c1e656d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=96=B9=E6=99=93=E5=85=B5?= Date: Mon, 3 Aug 2026 22:50:39 +0800 Subject: [PATCH 2/3] [server] Clean expired active segments without remote log - Roll expired non-empty active segments after all preceding segments pass retention checks and the high watermark reaches LEO. - Preserve two-pass retention and cover true remote-disabled and empty-active behavior. --- .../apache/fluss/server/log/LogTablet.java | 33 ++++---- .../fluss/server/log/LocalSegmentTTLTest.java | 83 +++++++++++++++++++ .../log/remote/TieredLocalSegmentTTLTest.java | 31 ++----- 3 files changed, 102 insertions(+), 45 deletions(-) create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/log/LocalSegmentTTLTest.java diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index f097165924..dd8efafb90 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -109,7 +109,6 @@ public final class LogTablet { private final Clock clock; private final boolean isChangeLog; private final long logTtlMs; - private final boolean remoteLogEnabled; @GuardedBy("lock") private volatile LogOffsetMetadata highWatermarkMetadata; @@ -161,8 +160,6 @@ private LogTablet( this.writerStateManager = writerStateManager; this.highWatermarkMetadata = new LogOffsetMetadata(0L); this.logTtlMs = logTtlMs; - this.remoteLogEnabled = - conf.get(ConfigOptions.REMOTE_LOG_TASK_INTERVAL_DURATION).toMillis() > 0L; this.scheduler = scheduler; // scheduler the writer expiration interval check. @@ -1362,30 +1359,28 @@ private List deletableExpiredSegments(long endOffset) throws IOExcep List deletableSegments = new ArrayList<>(); List logSegments = localLog.getSegments().values(); long now = clock.milliseconds(); - boolean shouldRoll = false; - for (int i = 0; i < logSegments.size(); i++) { - LogSegment segment = logSegments.get(i); - boolean active = i == logSegments.size() - 1; - if ((!active && logSegments.get(i + 1).getBaseOffset() > endOffset) - || !isSegmentExpired(now, segment, logTtlMs)) { + for (int i = 0; i < logSegments.size() - 1; i++) { + if (logSegments.get(i + 1).getBaseOffset() > endOffset + || !isSegmentExpired(now, logSegments.get(i), logTtlMs)) { break; } + deletableSegments.add(logSegments.get(i)); + } - if (active) { - shouldRoll = - remoteLogEnabled - && segment.getSizeInBytes() > 0 - && getHighWatermark() >= localLogEndOffset(); - break; - } - deletableSegments.add(segment); + if (deletableSegments.size() == logSegments.size() - 1) { + maybeRollExpiredActiveSegment(now, logSegments.get(logSegments.size() - 1)); } + return deletableSegments; + } - if (shouldRoll) { + private void maybeRollExpiredActiveSegment(long now, LogSegment activeSegment) + throws IOException { + if (activeSegment.getSizeInBytes() > 0 + && isSegmentExpired(now, activeSegment, logTtlMs) + && getHighWatermark() >= localLogEndOffset()) { roll(Optional.empty()); } - return deletableSegments; } @FunctionalInterface diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LocalSegmentTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LocalSegmentTTLTest.java new file mode 100644 index 0000000000..427ad90f5e --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LocalSegmentTTLTest.java @@ -0,0 +1,83 @@ +/* + * 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.server.log; + +import org.apache.fluss.config.ConfigOptions; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.server.replica.ReplicaTestBase; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.time.Duration; +import java.util.Collections; + +import static org.apache.fluss.record.TestData.DATA1_SCHEMA; +import static org.apache.fluss.record.TestData.DATA1_TABLE_ID; +import static org.apache.fluss.record.TestData.DATA1_TABLE_PATH; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests TTL-based cleanup when remote log is disabled. */ +final class LocalSegmentTTLTest extends ReplicaTestBase { + + @BeforeEach + public void setup() throws Exception { + super.setup(); + registerTableInZkClient( + DATA1_TABLE_PATH, + DATA1_SCHEMA, + DATA1_TABLE_ID, + Collections.emptyList(), + Collections.singletonMap(ConfigOptions.TABLE_LOG_TTL.key(), "1h")); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testExpiredActiveSegmentCleaned(boolean partitionTable) throws Exception { + TableBucket tb = + partitionTable + ? new TableBucket(DATA1_TABLE_ID, 0L, 0) + : new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(tb, partitionTable); + LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); + + assertThatThrownBy(() -> remoteLogManager.remoteLogTablet(tb)) + .isInstanceOf(IllegalStateException.class); + addMultiSegmentsToLogTablet(logTablet, 1); + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(1); + assertThat(logTablet.localLogStartOffset()).isEqualTo(0L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(0L); + + manualClock.advanceTime(Duration.ofHours(2)); + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(2); + assertThat(logTablet.localLogStartOffset()).isEqualTo(0L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(10L); + + logManager.cleanupExpiredLocalLogSegments(); + + assertThat(logTablet.getSegments()).hasSize(1); + assertThat(logTablet.localLogStartOffset()).isEqualTo(10L); + assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(10L); + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java index 3aee9e64f3..15328bc6e9 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java @@ -19,6 +19,7 @@ import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.server.log.LogSegment; import org.apache.fluss.server.log.LogTablet; import org.junit.jupiter.api.BeforeEach; @@ -77,17 +78,16 @@ void testInactiveTieredLocalSegmentRemovedAfterTtl(boolean partitionTable) throw assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); - // An empty active segment must not be rolled again while the previous segment is waiting - // for remote upload. - logManager.cleanupExpiredLocalLogSegments(); - assertThat(logTablet.getSegments()).hasSize(2); - remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); assertThat(remoteLog.getRemoteLogEndOffset()).hasValue(50L); logManager.cleanupExpiredLocalLogSegments(); assertThat(logTablet.getSegments()).hasSize(1); assertThat(logTablet.localLogStartOffset()).isEqualTo(50L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); + + LogSegment emptyActiveSegment = logTablet.activeLogSegment(); + logManager.cleanupExpiredLocalLogSegments(); + assertThat(logTablet.activeLogSegment()).isSameAs(emptyActiveSegment); } @ParameterizedTest @@ -120,27 +120,6 @@ void testExpiredActiveSegmentWaitsForHighWatermarkWithoutRemoteLogEndOffset( assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(50L); } - @ParameterizedTest - @ValueSource(booleans = {true, false}) - void testExpiredActiveSegmentNotRolledWhenRemoteLogDisabled(boolean partitionTable) - throws Exception { - conf.set(ConfigOptions.REMOTE_LOG_TASK_INTERVAL_DURATION, Duration.ZERO); - TableBucket tb = - partitionTable - ? new TableBucket(DATA1_TABLE_ID, 0L, 0) - : new TableBucket(DATA1_TABLE_ID, 0); - makeLogTableAsLeader(tb, partitionTable); - LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); - - addMultiSegmentsToLogTablet(logTablet, 5); - manualClock.advanceTime(Duration.ofHours(2)); - logManager.cleanupExpiredLocalLogSegments(); - - assertThat(logTablet.getSegments()).hasSize(1); - assertThat(logTablet.localLogStartOffset()).isEqualTo(40L); - assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); - } - @ParameterizedTest @ValueSource(booleans = {true, false}) void testExpiredActiveSegmentNotRolledWhenTtlDisabled(boolean partitionTable) throws Exception { From 57ebdeb18131f9b9b2b9f18e6f9538086236ac75 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=96=B9=E6=99=93=E5=85=B5?= Date: Wed, 5 Aug 2026 14:54:11 +0800 Subject: [PATCH 3/3] [server] Clarify active segment retention roll - Document the retention finder side effect.\n- State the conditions required to roll the active segment. --- .../main/java/org/apache/fluss/server/log/LogTablet.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index dd8efafb90..4d81d7c64f 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -1350,7 +1350,12 @@ private List deletableRemoteSegments(long endOffset) { return deletableSegments; } - /** Returns the contiguous prefix of expired segments and rolls an expired active segment. */ + /** + * Returns the contiguous prefix of expired inactive segments. + * + *

If all inactive segments are deletable, this also checks whether the active segment should + * be rolled. + */ private List deletableExpiredSegments(long endOffset) throws IOException { if (localLog.getSegments().isEmpty()) { return Collections.emptyList(); @@ -1374,6 +1379,7 @@ private List deletableExpiredSegments(long endOffset) throws IOExcep return deletableSegments; } + /** Rolls the active segment when it is non-empty, expired, and fully committed. */ private void maybeRollExpiredActiveSegment(long now, LogSegment activeSegment) throws IOException { if (activeSegment.getSizeInBytes() > 0