Mohsen Rezaei created FLINK-40931:
-------------------------------------
Summary: ForSt restore with object-store in restore mode fails
when scaling in
Key: FLINK-40931
URL: https://issues.apache.org/jira/browse/FLINK-40931
Project: Flink
Issue Type: Bug
Components: Runtime / State Backends
Affects Versions: 2.1.3, 2.2.1, 2.3.0, 2.0.2, 2.4.0
Reporter: Mohsen Rezaei
Related to [FLINK-38567|https://issues.apache.org/jira/browse/FLINK-38567]
While testing in-place scale-in of a query using the disaggregated state
backend with ForSt DB, I ran into the following exception when the job was
being restored at a lower parallelism (e.g. 4 -> 2):
{code:java}
org.apache.flink.runtime.state.BackendBuildingException: Caught unexpected
exception.
at
org.apache.flink.state.forst.ForStKeyedStateBackendBuilder.build(ForStKeyedStateBackendBuilder.java:315)
at
org.apache.flink.state.forst.ForStStateBackend.createAsyncKeyedStateBackend(ForStStateBackend.java:473)
Caused by: org.forstdb.RocksDBException: NotFound
at org.forstdb.Checkpoint.exportColumnFamily(Native Method)
at org.forstdb.Checkpoint.exportColumnFamily(Checkpoint.java:56)
at
org.apache.flink.state.forst.restore.ForStIncrementalRestoreOperation.exportColumnFamilies(ForStIncrementalRestoreOperation.java:761)
at
org.apache.flink.state.forst.restore.ForStIncrementalRestoreOperation.exportColumnFamiliesWithSstDataInKeyGroupsRange(ForStIncrementalRestoreOperation.java:900)
at
org.apache.flink.state.forst.restore.ForStIncrementalRestoreOperation.mergeStateHandlesWithClipAndIngest(ForStIncrementalRestoreOperation.java:687)
at
org.apache.flink.state.forst.restore.ForStIncrementalRestoreOperation.restoreFromMultipleStateHandles(ForStIncrementalRestoreOperation.java:362)
{code}
The Flink SQL query I used is a simple {{GROUP BY}} operation, started at
parallelism 4 and scaled in to 2:
{code:sql}
SELECT COUNT(page) as page_count, userid FROM user_audit GROUP BY userid;
{code}
This is the relevant configuration:
{code}
state.backend.type: forst
table.exec.async-state.enabled: true
execution.checkpointing.incremental: true
state.backend.forst.primary-dir: checkpoint-dir (s3p://...)
state.backend.forst.use-ingest-db-restore-mode: true
state.backend.forst.restore-overlap-fraction-threshold: 0.5
{code}
What I observed:
* The restore fails the same way on every retry, so the job ends up in
{{FAILED}} state after all retries are exhausted, if applicable.
* With {{state.backend.forst.use-ingest-db-restore-mode}} set to {{false}},
the same snapshot restores fine.
* Scale-out and same-parallelism restores are not affected.
* It only fails when a state handle lies entirely inside the new subtask's
key-group range. A scale-in shortly after a previous rescale may not hit it,
because those handles still hold data outside their range and are copied
instead.
It also reproduces without S3 in
{{ForStStateBackendV2Test#testAsyncStateBackendScaleDown}}, with the primary
directory on a {{LocalFileSystem}} subclass that mimics object-store directory
semantics: {{mkdirs()}} is a no-op, and a directory exists only while files
exist under it.
That is the behavior of the Presto S3 file system, and also of the native S3
file system.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)