[ 
https://issues.apache.org/jira/browse/FLINK-39412?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18125453#comment-18125453
 ] 

wind.wang commented on FLINK-39412:
-----------------------------------

We are hitting this in production with Flink CDC 3.6.0 (MySQL -> Paimon 
pipeline, Flink 1.20.5, application mode on Kubernetes, schema.change.behavior 
= TRY_EVOLVE).

*Symptom:* when a job is restarted (failover or restore from 
checkpoint/savepoint) while an ADD COLUMN DDL is still inside the binlog range 
being replayed, the same column is appended to the evolved schema again on 
every restart (column count 16 -> 17 -> 18 ...), until the job fails with 
{{java.lang.IllegalStateException: Duplicate key ...}} (same stack as 
FLINK-37537) and cannot recover without manual intervention.

*Local verification against release-3.6.0 SchemaUtils* (base schema col0..col3):
* Original 3.6.0: applying the same AddColumnEvent(col4) repeatedly (LAST / 
FIRST / AFTER / BEFORE) appends col4 every time (5 -> 6 -> 7 -> 8 columns) 
without any error; replaying "col4 AFTER col3" after col3 was dropped throws 
IllegalArgumentException.
* With the change from PR [#4370|https://github.com/apache/flink-cdc/pull/4370] 
applied on top of release-3.6.0 (the 10 added lines apply cleanly): column 
count stays at 5 in all cases, no exception. Only applyAddColumnEvent bytecode 
differs; public signatures unchanged.
* Replaying DropColumnEvent / RenameColumnEvent / AlterColumnTypeEvent for 
already-applied changes is a no-op in 3.6.0 both with and without the patch.

Since regular SchemaOperator, Pre/PostTransformOperator and SchemaManager all 
call SchemaUtils.applySchemaChangeEvent directly, making applyAddColumnEvent 
idempotent covers FLINK-37537 / FLINK-38830 / this issue; FLINK-37710 explains 
why the replay happens in the first place.

PR #4370 was auto-closed as stale on 2026-10-08 without review. We have asked 
the author to reopen it; if there is no response we would like to continue it 
in a new PR (crediting the original author). Any guidance from committers on 
the preferred approach would be appreciated.


> AddColumnEvent fails with duplicate field names when schema change events are 
> replayed after failover
> -----------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-39412
>                 URL: https://issues.apache.org/jira/browse/FLINK-39412
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>            Reporter: kianwee
>            Priority: Major
>              Labels: pull-request-available
>
> h3. Problem               
>           
>   When a Flink CDC pipeline recovers from a checkpoint/savepoint, binlog 
> events may be replayed, causing {{AddColumnEvent}} to be applied for columns 
> that already exist in the cached schema. This leads to a {{RowType}} 
> validation failure:       
>                             
>   {code:java}                                                                 
>                                                                               
>                                                                               
>            
>   
> org.apache.flink.cdc.runtime.operators.transform.exceptions.TransformException:
>  Failed to pre-transform with
>       AddColumnEvent{tableId=ecrm_btwl.kd_store_coupon, 
> addedColumns=[ColumnWithPosition{column=`valid_date` STRING, position=LAST, 
> existedColumnName=null}]}                                                     
>                                    
>   ...                                                                         
>                                                                               
>                                                                               
>            
>   Caused by: java.lang.IllegalArgumentException: Field names must be unique. 
> Found duplicates: [valid_date]                                                
>                                                                               
>             
>       at 
> org.apache.flink.cdc.common.types.RowType.validateFields(RowType.java:158)    
>                                                                               
>                                                                               
>   
>       at 
> org.apache.flink.cdc.runtime.operators.transform.PreTransformOperator.processElement(PreTransformOperator.java:230)
>   {code}                                                                      
>                                                                               
>                                                                               
>            
>                                                                               
>                         
>   h3. Root Cause                                                              
>                                                                               
>                                                                               
>            
>                                                                               
>                         
>   {{SchemaUtils.applyAddColumnEvent()}} blindly adds columns without checking 
> if a column with the same name already exists. While 
> {{isSchemaChangeEventRedundant()}} exists as a utility method, 
> {{PreTransformOperator.cacheChangeSchema()}} does  
>   not call it before applying schema changes.
>                                                                               
>                                                                               
>                                                                               
>            
>   This can be triggered when:                                                 
>                         
>   * A job restores from checkpoint/savepoint and the binlog offset rolls 
> back, replaying a historical {{ALTER TABLE ADD COLUMN}} DDL.
>   * The snapshot phase captures a schema that already includes the column, 
> but the binlog stream still contains the corresponding DDL event.             
>                                                                               
>               
>                                                                               
>                                                                               
>                                                                               
>            
>   h3. Fix                                                                     
>                                                                               
>                                                                               
>            
>                                                                               
>                                                                               
>                                                                               
>            
>   Add an idempotency check in {{SchemaUtils.applyAddColumnEvent()}} to skip 
> columns whose name already exists in the current schema. This is the most 
> defensive fix location since it protects all callers of 
> {{applySchemaChangeEvent()}}, not just 
>   {{PreTransformOperator}}.
>                                                                               
>                                                                               
>                                                                               
>            
>   PR: https://github.com/apache/flink-cdc/pull/4370                           
>                         
>                             
>   Priority: Major                                                             
>                                                                               
>                                                                               
>            
>                                     
>   Type: Bug    



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

Reply via email to