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

Yanquan Lv resolved FLINK-40552.
--------------------------------
    Fix Version/s: cdc-3.7.0
         Assignee: Mao Jiayi
       Resolution: Fixed

Merged in master via 7d560004bcbbd644ff2411b693e4f87f1646080a.

> Support ordered asynchronous execution for YAML Transform
> ---------------------------------------------------------
>
>                 Key: FLINK-40552
>                 URL: https://issues.apache.org/jira/browse/FLINK-40552
>             Project: Flink
>          Issue Type: Improvement
>          Components: Flink CDC
>            Reporter: Mao Jiayi
>            Assignee: Mao Jiayi
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: cdc-3.7.0
>
>
> Currently, the YAML PostTransform operator processes DataChangeEvents 
> synchronously.
> Transform expressions may invoke external AI model services or blocking UDFs. 
> In such cases, every record has to wait for the previous invocation to 
> finish, which limits pipeline throughput.
> This issue proposes an optional ordered asynchronous execution mode for 
> PostTransform.
> The implementation will:
> * Extract the common PostTransform processing logic into a processor shared 
> by synchronous and asynchronous execution.
> * Use AsyncDataStream.orderedWait for asynchronous execution.
> * Process DataChangeEvents concurrently while preserving their input order.
> * Treat SchemaChangeEvents as barriers:
>   ** Wait for all preceding DataChangeEvents to finish.
>   ** Apply the schema change.
>   ** Process subsequent DataChangeEvents with the updated schema.
> * Add the following pipeline options:
>   ** transform.async-execution.enabled
>   ** transform.async-execution.timeout
>   ** transform.async-execution.capacity
>   ** transform.async-execution.worker-threads
> * Preserve the current synchronous behavior when asynchronous execution is 
> disabled.
> * Persist the PostTransform schema view for checkpoint and savepoint recovery.
> * Re-emit the latest CreateTableEvent once per table before replayed 
> DataChangeEvents after recovery.
> * Propagate timeout and transform execution exceptions to the Flink job.
> The initial implementation will support checkpoint recovery and savepoint 
> recovery with unchanged parallelism. Restoring an asynchronous transform from 
> a savepoint with different parallelism is out of scope.



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

Reply via email to