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

Mohsen Rezaei updated FLINK-40931:
----------------------------------
    Description: 
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. On S3-like file systems (Presto, and also the native S3 file 
system, where {{mkdirs()}} is a no-op and a "directory" exists only if objects 
exist under the prefix) the renamed directory is reported as missing, and 
{{exportColumnFamily}} fails with {{NotFound}}.



That is the behavior of the Presto S3 file system, and also of the native S3 
file system.

  was:
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.


> 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.0.2, 2.3.0, 2.2.1, 2.1.3, 2.4.0
>            Reporter: Mohsen Rezaei
>            Priority: Critical
>
> 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. On S3-like file systems (Presto, and also the 
> native S3 file system, where {{mkdirs()}} is a no-op and a "directory" exists 
> only if objects exist under the prefix) the renamed directory is reported as 
> missing, and {{exportColumnFamily}} fails with {{NotFound}}.
> 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)

Reply via email to