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

ASF GitHub Bot updated FLINK-40043:
-----------------------------------
    Labels: pull-request-available  (was: )

> Paimon pipeline job couldn’t increase job parallelism
> -----------------------------------------------------
>
>                 Key: FLINK-40043
>                 URL: https://issues.apache.org/jira/browse/FLINK-40043
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.6.0
>         Environment: flink 1.20.3
> flinkcdc 3.6.0
>            Reporter: peiyu
>            Priority: Major
>              Labels: pull-request-available
>
> Fixed the issue where Paimon pipeline job couldn’t increase job parallelism.
> Reproduction method:
> 1. Stop the CDC job using a savepoint
> 2. Increase the parallelism of the task
> 3. Start the job using a savepoint
> Error message:
> {code:java}
> java.util.NoSuchElementException
>       at java.base/java.util.ArrayList$Itr.next(ArrayList.java:1000)
>       at 
> org.apache.flink.cdc.connectors.paimon.sink.v2.PaimonSink.restoreWriter(PaimonSink.java:101)
>       at 
> org.apache.flink.streaming.runtime.operators.sink.StatefulSinkWriterStateHandler.createWriter(StatefulSinkWriterStateHandler.java:120)
>       at 
> org.apache.flink.streaming.runtime.operators.sink.SinkWriterOperator.initializeState(SinkWriterOperator.java:153)
>       at 
> org.apache.flink.cdc.runtime.operators.sink.DataSinkWriterOperator.initializeState(DataSinkWriterOperator.java:132)
>       at 
> org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.initializeOperatorState(StreamOperatorStateHandler.java:147)
>       at 
> org.apache.flink.streaming.api.operators.AbstractStreamOperator.initializeState(AbstractStreamOperator.java:294)
>       at 
> org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.initializeStateAndOpenOperators(RegularOperatorChain.java:106)
>       at 
> org.apache.flink.streaming.runtime.tasks.StreamTask.restoreStateAndGates(StreamTask.java:858)
>       at 
> org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$restoreInternal$5(StreamTask.java:812)
>       at 
> org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.call(StreamTaskActionExecutor.java:55)
>       at 
> org.apache.flink.streaming.runtime.tasks.StreamTask.restoreInternal(StreamTask.java:812)
>       at 
> org.apache.flink.streaming.runtime.tasks.StreamTask.restore(StreamTask.java:771)
>       at 
> org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:970)
>       at 
> org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:939)
>       at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:763)
>       at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575)
>       at java.base/java.lang.Thread.run(Thread.java:829)
> {code}



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

Reply via email to