From c19c0cf6db9604cfb35d456c941023c20f48edc8 Mon Sep 17 00:00:00 2001 From: Seungjoo Choi Date: Thu, 10 Sep 2026 08:50:43 +0900 Subject: [PATCH 1/3] [FLINK-40354][state/forst] Configure the TTL compaction filter for State V2 states On the State V2 (async) path of the ForSt state backend the TTL compaction filter factory is attached to the column family when a TTL state is registered (ForStDBTtlCompactFiltersManager#setAndRegisterCompactFilterIfStateTtlV2), but the filter never receives its configuration: configCompactFilter() only accepts the V1 StateDescriptor and is only called by the sync backends. An unconfigured FlinkCompactionFilter stays disabled and keeps every entry, so expired state is never physically removed and state size grows without bound despite StateTtlConfig. This adds a configCompactFilter overload for the V2 StateDescriptor (sharing the existing body) and calls it from ForStKeyedStateBackend#createStateInternal right after the state is registered, mirroring ForStSyncKeyedStateBackend and RocksDBKeyedStateBackend. A compactState() test hook is added to the async backend as well, and ForStTtlCompactFilterTest verifies that expired entries of a V2 ValueState are dropped by compaction. Co-Authored-By: Claude Code --- .../ForStDBTtlCompactFiltersManager.java | 43 ++++- .../state/forst/ForStKeyedStateBackend.java | 16 ++ .../flink/state/forst/ForStTestUtils.java | 13 +- .../forst/ForStTtlCompactFilterTest.java | 168 ++++++++++++++++++ 4 files changed, 233 insertions(+), 7 deletions(-) create mode 100644 flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java diff --git a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java index 89b55b9691874..d4748c72eb168 100644 --- a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java +++ b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java @@ -149,15 +149,46 @@ protected void log(InfoLogLevel infoLogLevel, String logMsg) { public void configCompactFilter( @Nonnull StateDescriptor stateDesc, TypeSerializer stateSerializer) { - StateTtlConfig ttlConfig = stateDesc.getTtlConfig(); + configCompactFilter( + stateDesc.getName(), + stateDesc.getTtlConfig(), + stateDesc instanceof ListStateDescriptor, + stateDesc instanceof MapStateDescriptor, + stateSerializer); + } + + /** + * Configures the TTL compaction filter for a state created through the State V2 API. This is + * the counterpart of {@link #configCompactFilter(StateDescriptor, TypeSerializer)}: registering + * the filter factory on the column family alone (see {@link + * #setAndRegisterCompactFilterIfStateTtlV2}) is not sufficient, the native filter stays + * disabled until it receives its {@link FlinkCompactionFilter.Config}. + */ + public void configCompactFilter( + @Nonnull org.apache.flink.api.common.state.v2.StateDescriptor stateDesc, + TypeSerializer stateSerializer) { + org.apache.flink.api.common.state.v2.StateDescriptor.Type type = stateDesc.getType(); + configCompactFilter( + stateDesc.getStateId(), + stateDesc.getTtlConfig(), + type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.LIST, + type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.MAP, + stateSerializer); + } + + private void configCompactFilter( + String stateName, + StateTtlConfig ttlConfig, + boolean isListState, + boolean isMapState, + TypeSerializer stateSerializer) { if (ttlConfig.isEnabled() && ttlConfig.getCleanupStrategies().inRocksdbCompactFilter()) { FlinkCompactionFilterFactory compactionFilterFactory = - compactionFilterFactories.get(stateDesc.getName()); + compactionFilterFactories.get(stateName); Preconditions.checkNotNull(compactionFilterFactory); long ttl = ttlConfig.getTimeToLive().toMillis(); - ColumnFamilyOptions columnFamilyOptions = - columnFamilyOptionsMap.get(stateDesc.getName()); + ColumnFamilyOptions columnFamilyOptions = columnFamilyOptionsMap.get(stateName); Preconditions.checkNotNull(columnFamilyOptions); StateTtlConfig.RocksdbCompactFilterCleanupStrategy rocksdbCompactFilterCleanupStrategy = @@ -178,7 +209,7 @@ public void configCompactFilter( } FlinkCompactionFilter.Config config; - if (stateDesc instanceof ListStateDescriptor) { + if (isListState) { TypeSerializer elemSerializer = ((ListSerializer) stateSerializer).getElementSerializer(); int len = elemSerializer.getLength(); @@ -195,7 +226,7 @@ public void configCompactFilter( queryTimeAfterNumEntries, new ListElementFilterFactory<>(elemSerializer.duplicate())); } - } else if (stateDesc instanceof MapStateDescriptor) { + } else if (isMapState) { config = FlinkCompactionFilter.Config.createForMap(ttl, queryTimeAfterNumEntries); } else { config = FlinkCompactionFilter.Config.createForValue(ttl, queryTimeAfterNumEntries); diff --git a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStKeyedStateBackend.java b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStKeyedStateBackend.java index 8160252fb73c1..fff4342568ace 100644 --- a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStKeyedStateBackend.java +++ b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStKeyedStateBackend.java @@ -66,6 +66,7 @@ import org.forstdb.ColumnFamilyHandle; import org.forstdb.ColumnFamilyOptions; import org.forstdb.RocksDB; +import org.forstdb.RocksDBException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -291,6 +292,14 @@ public S createStateInternal( Tuple2> registerResult = tryRegisterKvStateInformation(stateDesc, namespaceSerializer); + // The compaction filter factory is attached to the column family when the state is + // registered, but the native filter only becomes active once it receives the TTL + // configuration. Mirror what the sync backend does after registration. + if (ttlCompactFiltersManager != null) { + ttlCompactFiltersManager.configCompactFilter( + stateDesc, registerResult.f1.getStateSerializer()); + } + ColumnFamilyHandle columnFamilyHandle = registerResult.f0; switch (stateDesc.getType()) { @@ -624,6 +633,13 @@ public boolean isSafeToReuseKVState() { return true; } + @VisibleForTesting + public void compactState(StateDescriptor stateDesc) throws RocksDBException { + ForStOperationUtils.ForStKvStateInfo kvStateInfo = + kvStateInformation.get(stateDesc.getStateId()); + db.compactRange(kvStateInfo.columnFamilyHandle); + } + @VisibleForTesting Path getLocalBasePath() { return optionsContainer.getPathContainer().getLocalBasePath(); diff --git a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTestUtils.java b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTestUtils.java index 89e31baee0865..7425c7e450be1 100644 --- a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTestUtils.java +++ b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTestUtils.java @@ -41,6 +41,17 @@ public static ForStKeyedStateBackend createKeyedStateBackend( TypeSerializer keySerializer, Collection stateHandles) throws IOException { + return createKeyedStateBackend( + forStStateBackend, env, keySerializer, stateHandles, TtlTimeProvider.DEFAULT); + } + + public static ForStKeyedStateBackend createKeyedStateBackend( + ForStStateBackend forStStateBackend, + Environment env, + TypeSerializer keySerializer, + Collection stateHandles, + TtlTimeProvider ttlTimeProvider) + throws IOException { return forStStateBackend.createAsyncKeyedStateBackend( new KeyedStateBackendParametersImpl<>( @@ -51,7 +62,7 @@ public static ForStKeyedStateBackend createKeyedStateBackend( 1, new KeyGroupRange(0, 0), env.getTaskKvStateRegistry(), - TtlTimeProvider.DEFAULT, + ttlTimeProvider, new UnregisteredMetricsGroup(), (name, value) -> {}, stateHandles, diff --git a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java new file mode 100644 index 0000000000000..601c8cacc240d --- /dev/null +++ b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java @@ -0,0 +1,168 @@ +/* + * 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.flink.state.forst; + +import org.apache.flink.api.common.operators.MailboxExecutor; +import org.apache.flink.api.common.state.StateTtlConfig; +import org.apache.flink.api.common.state.v2.ValueStateDescriptor; +import org.apache.flink.api.common.typeutils.base.StringSerializer; +import org.apache.flink.configuration.Configuration; +import org.apache.flink.core.fs.FileSystem; +import org.apache.flink.runtime.asyncprocessing.EpochManager; +import org.apache.flink.runtime.asyncprocessing.RecordContext; +import org.apache.flink.runtime.asyncprocessing.StateExecutionController; +import org.apache.flink.runtime.asyncprocessing.declare.DeclarationManager; +import org.apache.flink.runtime.operators.testutils.MockEnvironment; +import org.apache.flink.runtime.state.VoidNamespace; +import org.apache.flink.runtime.state.VoidNamespaceSerializer; +import org.apache.flink.runtime.state.ttl.TtlTimeProvider; +import org.apache.flink.runtime.state.v2.internal.InternalValueState; +import org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor; +import org.apache.flink.streaming.runtime.tasks.mailbox.MailboxExecutorImpl; +import org.apache.flink.streaming.runtime.tasks.mailbox.TaskMailboxImpl; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.File; +import java.time.Duration; +import java.util.Collections; +import java.util.concurrent.atomic.AtomicLong; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests that the TTL compaction filter is actually configured (and therefore removes expired + * entries during compaction) for states created through the State V2 API on {@link + * ForStKeyedStateBackend}. + */ +class ForStTtlCompactFilterTest { + + private static final Duration TTL = Duration.ofMillis(1000); + + private final AtomicLong currentTime = new AtomicLong(0L); + + private ForStKeyedStateBackend keyedBackend; + private StateExecutionController aec; + private RecordContext context; + private MockEnvironment env; + + @BeforeEach + void setup(@TempDir File temporaryFolder) throws Exception { + FileSystem.initialize(new Configuration(), null); + Configuration configuration = new Configuration(); + configuration.set(ForStOptions.PRIMARY_DIRECTORY, temporaryFolder.toURI().toString()); + ForStStateBackend forStStateBackend = + new ForStStateBackend().configure(configuration, null); + + env = ForStStateTestBase.getMockEnvironment(temporaryFolder); + + TtlTimeProvider timeProvider = currentTime::get; + keyedBackend = + ForStTestUtils.createKeyedStateBackend( + forStStateBackend, + env, + StringSerializer.INSTANCE, + Collections.emptyList(), + timeProvider); + + MailboxExecutor mailboxExecutor = + new MailboxExecutorImpl( + new TaskMailboxImpl(), 0, StreamTaskActionExecutor.IMMEDIATE); + aec = + new StateExecutionController<>( + mailboxExecutor, + (a, b) -> {}, + keyedBackend.createStateExecutor(), + new DeclarationManager(), + EpochManager.ParallelMode.SERIAL_BETWEEN_EPOCH, + 1, + 100, + 0, + 1, + null, + null); + keyedBackend.setup(aec); + } + + @AfterEach + void tearDown() throws Exception { + keyedBackend.close(); + env.close(); + } + + @Test + void testExpiredEntriesAreRemovedByCompaction() throws Exception { + // ReturnExpiredIfNotCleanedUp: the read path does not hide expired entries, so a null + // value after compaction proves that the compaction filter removed the entry physically. + StateTtlConfig ttlConfig = + StateTtlConfig.newBuilder(TTL) + .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) + .setStateVisibility( + StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp) + .cleanupInRocksdbCompactFilter(1L) + .build(); + ValueStateDescriptor descriptor = + new ValueStateDescriptor<>("ttl-value-state", StringSerializer.INSTANCE); + descriptor.enableTimeToLive(ttlConfig); + + InternalValueState state = + keyedBackend.createState( + VoidNamespace.INSTANCE, VoidNamespaceSerializer.INSTANCE, descriptor); + + currentTime.set(0L); + for (String key : new String[] {"k1", "k2", "k3"}) { + setCurrentContext(key); + state.update("v-" + key); + drain(); + } + + // Move the clock past the TTL and write one fresh entry: the old ones are expired, the + // new one is not. + currentTime.set(TTL.toMillis() + 1); + setCurrentContext("k4"); + state.update("v-k4"); + drain(); + + // Without a configured compaction filter this is a plain rewrite and nothing is dropped. + keyedBackend.compactState(descriptor); + + for (String key : new String[] {"k1", "k2", "k3"}) { + setCurrentContext(key); + assertThat(state.value()).as("expired entry %s should be removed", key).isNull(); + drain(); + } + setCurrentContext("k4"); + assertThat(state.value()).isEqualTo("v-k4"); + drain(); + } + + private void setCurrentContext(String key) { + context = aec.buildContext(key, key); + context.retain(); + aec.setCurrentContext(context); + } + + private void drain() { + context.release(); + aec.drainInflightRecords(0); + } +} From b64175b936f221855f2ad4f2a4a85f86ea4e0937 Mon Sep 17 00:00:00 2001 From: Seungjoo Choi Date: Sat, 12 Sep 2026 00:47:47 +0900 Subject: [PATCH 2/3] [FLINK-40354][state/forst] Use the element serializer directly for V2 list states A V2 ListStateDescriptor carries the element serializer itself (for a TTL state a TtlSerializer), not a ListSerializer> as in V1, so unwrapping it as a ListSerializer threw a ClassCastException when a TTL-enabled V2 ListState was created. Unwrap only on the V1 path and add V2 ListState regression tests (variable- and fixed-length elements). --- .../ForStDBTtlCompactFiltersManager.java | 38 ++++++--- .../forst/ForStTtlCompactFilterTest.java | 83 ++++++++++++++++--- 2 files changed, 99 insertions(+), 22 deletions(-) diff --git a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java index d4748c72eb168..07972ac4faef9 100644 --- a/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java +++ b/flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/ForStDBTtlCompactFiltersManager.java @@ -44,6 +44,7 @@ import org.slf4j.LoggerFactory; import javax.annotation.Nonnull; +import javax.annotation.Nullable; import java.io.IOException; import java.time.Duration; @@ -149,12 +150,17 @@ protected void log(InfoLogLevel infoLogLevel, String logMsg) { public void configCompactFilter( @Nonnull StateDescriptor stateDesc, TypeSerializer stateSerializer) { + // A V1 list state is registered with a ListSerializer>; the compaction filter + // works on the elements, so unwrap it here. + TypeSerializer listElementSerializer = + stateDesc instanceof ListStateDescriptor + ? ((ListSerializer) stateSerializer).getElementSerializer() + : null; configCompactFilter( stateDesc.getName(), stateDesc.getTtlConfig(), - stateDesc instanceof ListStateDescriptor, - stateDesc instanceof MapStateDescriptor, - stateSerializer); + listElementSerializer, + stateDesc instanceof MapStateDescriptor); } /** @@ -168,20 +174,29 @@ public void configCompactFilter( @Nonnull org.apache.flink.api.common.state.v2.StateDescriptor stateDesc, TypeSerializer stateSerializer) { org.apache.flink.api.common.state.v2.StateDescriptor.Type type = stateDesc.getType(); + // Unlike V1, a V2 ListStateDescriptor carries the element serializer itself (for a TTL + // state: TtlSerializer), and ForStListState stores the elements delimited with that + // serializer. So the registered state serializer is already the list element serializer. + TypeSerializer listElementSerializer = + type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.LIST + ? stateSerializer + : null; configCompactFilter( stateDesc.getStateId(), stateDesc.getTtlConfig(), - type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.LIST, - type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.MAP, - stateSerializer); + listElementSerializer, + type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.MAP); } + /** + * @param listElementSerializer the element serializer of a list state (the TTL-wrapped element + * serializer), or {@code null} if the state is not a list state. + */ private void configCompactFilter( String stateName, StateTtlConfig ttlConfig, - boolean isListState, - boolean isMapState, - TypeSerializer stateSerializer) { + @Nullable TypeSerializer listElementSerializer, + boolean isMapState) { if (ttlConfig.isEnabled() && ttlConfig.getCleanupStrategies().inRocksdbCompactFilter()) { FlinkCompactionFilterFactory compactionFilterFactory = compactionFilterFactories.get(stateName); @@ -209,9 +224,8 @@ private void configCompactFilter( } FlinkCompactionFilter.Config config; - if (isListState) { - TypeSerializer elemSerializer = - ((ListSerializer) stateSerializer).getElementSerializer(); + if (listElementSerializer != null) { + TypeSerializer elemSerializer = listElementSerializer; int len = elemSerializer.getLength(); if (len > 0) { config = diff --git a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java index 601c8cacc240d..289466ffff4f3 100644 --- a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java +++ b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java @@ -20,7 +20,10 @@ import org.apache.flink.api.common.operators.MailboxExecutor; import org.apache.flink.api.common.state.StateTtlConfig; +import org.apache.flink.api.common.state.v2.ListStateDescriptor; import org.apache.flink.api.common.state.v2.ValueStateDescriptor; +import org.apache.flink.api.common.typeutils.TypeSerializer; +import org.apache.flink.api.common.typeutils.base.LongSerializer; import org.apache.flink.api.common.typeutils.base.StringSerializer; import org.apache.flink.configuration.Configuration; import org.apache.flink.core.fs.FileSystem; @@ -32,6 +35,7 @@ import org.apache.flink.runtime.state.VoidNamespace; import org.apache.flink.runtime.state.VoidNamespaceSerializer; import org.apache.flink.runtime.state.ttl.TtlTimeProvider; +import org.apache.flink.runtime.state.v2.internal.InternalListState; import org.apache.flink.runtime.state.v2.internal.InternalValueState; import org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor; import org.apache.flink.streaming.runtime.tasks.mailbox.MailboxExecutorImpl; @@ -44,6 +48,7 @@ import java.io.File; import java.time.Duration; +import java.util.Arrays; import java.util.Collections; import java.util.concurrent.atomic.AtomicLong; @@ -109,20 +114,23 @@ void tearDown() throws Exception { env.close(); } + /** + * ReturnExpiredIfNotCleanedUp: the read path does not hide expired entries, so an entry that is + * gone after compaction proves that the compaction filter removed it physically. + */ + private static StateTtlConfig ttlConfig() { + return StateTtlConfig.newBuilder(TTL) + .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) + .setStateVisibility(StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp) + .cleanupInRocksdbCompactFilter(1L) + .build(); + } + @Test void testExpiredEntriesAreRemovedByCompaction() throws Exception { - // ReturnExpiredIfNotCleanedUp: the read path does not hide expired entries, so a null - // value after compaction proves that the compaction filter removed the entry physically. - StateTtlConfig ttlConfig = - StateTtlConfig.newBuilder(TTL) - .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) - .setStateVisibility( - StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp) - .cleanupInRocksdbCompactFilter(1L) - .build(); ValueStateDescriptor descriptor = new ValueStateDescriptor<>("ttl-value-state", StringSerializer.INSTANCE); - descriptor.enableTimeToLive(ttlConfig); + descriptor.enableTimeToLive(ttlConfig()); InternalValueState state = keyedBackend.createState( @@ -155,6 +163,61 @@ void testExpiredEntriesAreRemovedByCompaction() throws Exception { drain(); } + /** + * A V2 list state is registered with the element serializer itself (a TtlSerializer), not a + * ListSerializer as in V1 — configuring the filter must not assume the V1 shape. + * Variable-length elements take the element-filter path of the native filter. + */ + @Test + void testExpiredListElementsAreRemovedByCompaction() throws Exception { + testExpiredListElementsAreRemovedByCompaction( + "ttl-list-state-var", StringSerializer.INSTANCE, "a", "b", "c"); + } + + /** Fixed-length elements take the fixed-element-length path of the native filter. */ + @Test + void testExpiredFixedLengthListElementsAreRemovedByCompaction() throws Exception { + testExpiredListElementsAreRemovedByCompaction( + "ttl-list-state-fixed", LongSerializer.INSTANCE, 1L, 2L, 3L); + } + + private void testExpiredListElementsAreRemovedByCompaction( + String stateName, TypeSerializer elementSerializer, E e1, E e2, E fresh) + throws Exception { + ListStateDescriptor descriptor = new ListStateDescriptor<>(stateName, elementSerializer); + descriptor.enableTimeToLive(ttlConfig()); + + InternalListState state = + keyedBackend.createState( + VoidNamespace.INSTANCE, VoidNamespaceSerializer.INSTANCE, descriptor); + + // k1: two elements that will expire, plus one fresh element added after the TTL. + // k2: only expired elements. + currentTime.set(0L); + setCurrentContext("k1"); + state.update(Arrays.asList(e1, e2)); + drain(); + setCurrentContext("k2"); + state.update(Arrays.asList(e1, e2)); + drain(); + + currentTime.set(TTL.toMillis() + 1); + setCurrentContext("k1"); + state.add(fresh); + drain(); + + keyedBackend.compactState(descriptor); + + setCurrentContext("k1"); + assertThat(state.get()) + .as("expired elements of k1 should be removed, the fresh one kept") + .containsExactly(fresh); + drain(); + setCurrentContext("k2"); + assertThat(state.get()).as("k2 should have no elements left").isEmpty(); + drain(); + } + private void setCurrentContext(String key) { context = aec.buildContext(key, key); context.retain(); From a729c70b538895cd046e65e50494b4f62d3d5814 Mon Sep 17 00:00:00 2001 From: Seungjoo Choi Date: Sat, 12 Sep 2026 01:01:23 +0900 Subject: [PATCH 3/3] [FLINK-40354][state/forst] Cover V2 MapState in ForStTtlCompactFilterTest --- .../forst/ForStTtlCompactFilterTest.java | 37 +++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java index 289466ffff4f3..1c5c29e81d577 100644 --- a/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java +++ b/flink-state-backends/flink-statebackend-forst/src/test/java/org/apache/flink/state/forst/ForStTtlCompactFilterTest.java @@ -21,6 +21,7 @@ import org.apache.flink.api.common.operators.MailboxExecutor; import org.apache.flink.api.common.state.StateTtlConfig; import org.apache.flink.api.common.state.v2.ListStateDescriptor; +import org.apache.flink.api.common.state.v2.MapStateDescriptor; import org.apache.flink.api.common.state.v2.ValueStateDescriptor; import org.apache.flink.api.common.typeutils.TypeSerializer; import org.apache.flink.api.common.typeutils.base.LongSerializer; @@ -36,6 +37,7 @@ import org.apache.flink.runtime.state.VoidNamespaceSerializer; import org.apache.flink.runtime.state.ttl.TtlTimeProvider; import org.apache.flink.runtime.state.v2.internal.InternalListState; +import org.apache.flink.runtime.state.v2.internal.InternalMapState; import org.apache.flink.runtime.state.v2.internal.InternalValueState; import org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor; import org.apache.flink.streaming.runtime.tasks.mailbox.MailboxExecutorImpl; @@ -181,6 +183,41 @@ void testExpiredFixedLengthListElementsAreRemovedByCompaction() throws Exception "ttl-list-state-fixed", LongSerializer.INSTANCE, 1L, 2L, 3L); } + /** + * Map entries carry their own TTL timestamp behind the null flag of the serialized user value; + * the filter is configured in map mode for that layout. + */ + @Test + void testExpiredMapEntriesAreRemovedByCompaction() throws Exception { + MapStateDescriptor descriptor = + new MapStateDescriptor<>( + "ttl-map-state", StringSerializer.INSTANCE, StringSerializer.INSTANCE); + descriptor.enableTimeToLive(ttlConfig()); + + InternalMapState state = + keyedBackend.createState( + VoidNamespace.INSTANCE, VoidNamespaceSerializer.INSTANCE, descriptor); + + currentTime.set(0L); + setCurrentContext("k1"); + state.put("uk1", "v1"); + state.put("uk2", "v2"); + drain(); + + currentTime.set(TTL.toMillis() + 1); + setCurrentContext("k1"); + state.put("uk3", "v3"); + drain(); + + keyedBackend.compactState(descriptor); + + setCurrentContext("k1"); + assertThat(state.get("uk1")).as("expired entry uk1 should be removed").isNull(); + assertThat(state.get("uk2")).as("expired entry uk2 should be removed").isNull(); + assertThat(state.get("uk3")).isEqualTo("v3"); + drain(); + } + private void testExpiredListElementsAreRemovedByCompaction( String stateName, TypeSerializer elementSerializer, E e1, E e2, E fresh) throws Exception {