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)