Francis created FLINK-40157:
-------------------------------

             Summary: MapState#asyncPutAll on ForSt: reads chained on the 
returned StateFuture see the new keys but null values
                 Key: FLINK-40157
                 URL: https://issues.apache.org/jira/browse/FLINK-40157
             Project: Flink
          Issue Type: Bug
          Components: Runtime / State Backends
    Affects Versions: 2.2.1
         Environment: Flink 2.2.0 using ForSt with a S3 remote backend. 
            Reporter: Francis


*Environment*
 * Flink 2.2.0
 * ForSt state backend with async state enabled (\{{enableAsyncState()}}, State 
V2 API)
 * \{{KeyedProcessFunction}} with a \{{MapState<String, String>}}

*Problem*

When writing to a \{{MapState}} with \{{asyncPutAll}} and chaining reads on the 
StateFuture it returns, the reads see all of the newly written keys but every 
value is null. Concretely, after the \{{asyncPutAll}} future completes:
 * \{{asyncKeys()}} returns all the keys
 * \{{asyncValues()}} returns nothing
 * \{{StateFutureUtils.toIterable(asyncEntries())}} produces entries where 
every value is null

The data is not lost. We took a checkpoint from the affected job and inspected 
it, and both the keys and the values are present in the MapState. So the write 
does reach the backend. It looks like a race where the future returned by 
\{{asyncPutAll}} completes before the values are visible to reads chained on it.

This only reproduces on ForSt. Our integration tests run the same pipeline on 
MiniCluster and do not see the issue. Running the official Flink docker image 
in single instance mode also does not reproduce it. It reproduces consistently 
in our cluster deployment on ForSt.

*Minimal reproduction (Kotlin)*

{code:kotlin}
class TagEnricher : KeyedProcessFunction<Long, Event, EnrichedEvent>() {

    @field:Transient
    private lateinit var tagState: MapState<String, String>

    override fun open(openContext: OpenContext?) {
        tagState = runtimeContext.getMapState(
            MapStateDescriptor(
                "TagEnricher-tagState",
                BasicTypeInfo.STRING_TYPE_INFO,
                BasicTypeInfo.STRING_TYPE_INFO,
            ),
        )
    }

    override fun processElement(
        event: Event,
        context: Context?,
        out: Collector<EnrichedEvent>,
    ) {
        tagState.asyncPutAll(event.tags)
            .thenCompose {
                StateFutureUtils.toIterable(tagState.asyncEntries())
            }
            .thenAccept { entries ->
                // On ForSt: entries contains every key from event.tags,
                // but every value is null.
                out.collect(EnrichedEvent(event.id, entries.toList()))
            }
    }
}
{code}

*Expected behavior*

Once the future returned by \{{asyncPutAll}} completes, reads chained on that 
future should see both the keys and the values that were written.

*Actual behavior*

Chained reads see the keys with null values, and \{{asyncValues()}} is empty. 
Checkpoints taken from the same job contain the full key/value pairs.

*Workaround*

Replacing the single \{{asyncPutAll}} with one \{{asyncPut}} per entry, 
combined with \{{StateFutureUtils.combineAll}}, makes the chained reads return 
correct values:

{code:kotlin}
val putFutures = event.tags.map { (key, value) ->
    tagState.asyncPut(key, value)
}
StateFutureUtils.combineAll(putFutures)
    .thenCompose {
        StateFutureUtils.toIterable(tagState.asyncEntries())
    }
{code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to