[ 
https://issues.apache.org/jira/browse/FLINK-40157?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Francis updated FLINK-40157:
----------------------------
    Description: 
*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}
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}
val putFutures = event.tags.map { (key, value) ->
    tagState.asyncPut(key, value)
}
StateFutureUtils.combineAll(putFutures)
    .thenCompose {
        StateFutureUtils.toIterable(tagState.asyncEntries())
    }
{code}

  was:
*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}
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}


> 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
>            Priority: Major
>
> *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}
> 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}
> 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