[
https://issues.apache.org/jira/browse/FLINK-40354?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40354:
-----------------------------------
Labels: pull-request-available (was: )
> StateTtlConfig background cleanup does not work for async State v2 MapState
> on ForSt state backend
> --------------------------------------------------------------------------------------------------
>
> Key: FLINK-40354
> URL: https://issues.apache.org/jira/browse/FLINK-40354
> Project: Flink
> Issue Type: Bug
> Affects Versions: 2.2.0
> Environment: Flink version: 2.2.0
> State backend: ForSt (via {{{}flink-statebackend-rocksdb{}}})
> State API: {{org.apache.flink.api.common.state.v2.MapState}} (async)
> Java: 21
> Checkpointing: incremental, interval 30s
> Reporter: Mykola Lukashchuk
> Priority: Major
> Labels: pull-request-available
>
> When using the new async State v2 API
> (org.apache.flink.api.common.state.v2.MapState obtained via
> org.apache.flink.api.common.state.v2.MapStateDescriptor) with StateTtlConfig
> enabled, expired state entries are never cleaned up in the background when
> running on the ForSt state backend. State size grows indefinitely despite TTL
> being properly configured.
> h3. Reproducer
> {code:java}
> import org.apache.flink.api.common.state.StateTtlConfig;
> import org.apache.flink.api.common.state.v2.MapState;
> import org.apache.flink.api.common.state.v2.MapStateDescriptor;
> import org.apache.flink.api.common.typeinfo.Types;
> import org.apache.flink.api.java.tuple.Tuple2;
> import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
> import java.time.Duration;
> public class DeduplicationProcessFunc extends KeyedProcessFunction<String,
> MyEvent, MyEvent> {
> private transient MapState<Tuple2<Long, Long>, Boolean> deduplicationState;
> @Override
> public void open(OpenContext openContext) throws Exception {
> var ttlConfig = StateTtlConfig
> .newBuilder(Duration.ofDays(7))
> .updateTtlOnCreateAndWrite()
> .returnExpiredIfNotCleanedUp()
> .useProcessingTime()
> .cleanupIncrementally(10, true)
> .build();
> var descriptor = new MapStateDescriptor<Tuple2<Long, Long>, Boolean>(
> "processed-events",
> Types.TUPLE(Types.LONG, Types.LONG),
> Types.BOOLEAN
> );
> descriptor.enableTimeToLive(ttlConfig);
> this.deduplicationState = getRuntimeContext().getMapState(descriptor);
> }
> @Override
> public void processElement(MyEvent value, Context ctx, Collector<MyEvent>
> out) throws Exception {
> Tuple2<Long, Long> hash = computeHash(value);
> deduplicationState.asyncContains(hash).thenAccept(exists -> {
> if (exists == null || !exists) {
> out.collect(value);
> deduplicationState.asyncPut(hash, Boolean.TRUE);
> }
> });
> }
> } {code}
> h3. Expected Behavior
> State entries older than 7 days should be cleaned up via:
> ForSt compaction filter (background cleanup), and/or
> Lazy expiration on read path (when {{asyncContains}} is called for an expired
> entry)
> Checkpoint size should stabilize after the state reaches a steady-state
> window matching the TTL duration.
> h3. Actual Behavior
> State size grows continuously and indefinitely. Entries are never removed
> regardless of their age. After running for 3+ weeks with a 7-day TTL, the
> state size is approximately 3x larger than expected maximum.
> Neither background cleanup nor read-path lazy expiration appears to function.
> h3. Analysis
> The Flink documentation states: "Currently, heap state backend relies on
> incremental cleanup and RocksDB backend uses compaction filter for background
> cleanup."
> Two separate issues appear to be present:
> {{cleanupIncrementally()}} is silently ignored on ForSt: This cleanup
> strategy only applies to the heap state backend. When used with ForSt, no
> error or warning is produced, and no cleanup occurs. The API allows
> configuring it without any indication that it will be ineffective.
> ForSt compaction filter may not be registered for State v2 path: The
> traditional RocksDB state backend uses {{FlinkCompactionFilterFactory}} to
> enable background TTL cleanup during compaction. The new ForSt async state
> creation path (serving
> {{{}org.apache.flink.api.common.state.v2.MapStateDescriptor{}}}) may not
> register this compaction filter for column families created through the v2
> API, meaning the default background cleanup (which should be active even
> without explicit cleanup strategy configuration) does not apply.
> h3. Relevant Code Paths to Investigate
> {{ForStKeyedStateBackend}} — how async v2 state instances are created and
> whether TTL wrapping is applied
> Whether {{ForStCompactionFilter}} or equivalent is registered for column
> families serving v2 state descriptors with TTL enabled
> Whether the TTL timestamp is stored alongside values when using the v2
> {{MapStateDescriptor}} with {{enableTimeToLive()}}
> {{AbstractMapState}} in {{org.apache.flink.runtime.state.v2}} — whether it
> integrates with TTL mechanisms
> h3. Workaround
> No effective workaround found. Possible mitigations:
> Manual cleanup via processing-time timers (defeats the purpose of
> deduplication state and adds significant complexity)
> Switching back to legacy State v1 API — not feasible since the DataStream API
> v2 runtime requires v2 state
--
This message was sent by Atlassian Jira
(v8.20.10#820010)