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

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

> Re-added table with drifted schema is silently skipped: downstream table not 
> re-created and stale schema state causes write failures
> ------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40795
>                 URL: https://issues.apache.org/jira/browse/FLINK-40795
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.6.0
>            Reporter: peiyu
>            Priority: Major
>              Labels: pull-request-available
>
> h2. *Symptom*
> In a pipeline with `scan.newly-added-table.enabled: true` (verified with 
> MySQL -> Paimon),
> when a table is removed from the pipeline configuration and later added back, 
> while its
> upstream schema drifted in between (e.g. a column was dropped while it was 
> out of capture):
> 1. The downstream table is NOT re-created automatically, even though the 
> source re-snapshots
>    the table and emits a `CreateTableEvent` carrying the latest schema.
> 2. If the downstream table is created manually with the current upstream 
> schema, the job
>    fails with metadata/schema mismatch errors when writing data, because the 
> pipeline keeps
>    coercing records into the stale schema recorded in state.
> h2. *Steps to Reproduce*
> 1. Run a MySQL -> Paimon pipeline with `scan.newly-added-table.enabled: true` 
> capturing
>    table `t`; take a savepoint and stop.
> 2. Remove `t` from the pipeline's `tables` configuration, and drop the Paimon 
> table `t`.
> 3. Restart from the savepoint, run a few checkpoints, stop with a savepoint 
> again.
> 4. Drop a column from `t` in MySQL (while it is out of capture), then add `t` 
> back to the
>    pipeline configuration.
> 5. Restart from the savepoint.
> Actual: the source re-snapshots `t`, but the Paimon table is never created. 
> The coordinator
> logs "Schema change event CreateTableEvent{...} is redundant for current 
> schema ..., just
> skip it."
> Expected: `t` is re-created downstream with its current schema and data flows 
> normally.
> The same reproduces at unit level on current master: restore a regular 
> `SchemaCoordinator`
> from a checkpoint containing schema V1 of table `t`, then send 
> `CreateTableEvent(t, V2)`
> with V2 != V1. `SchemaChangeResponse#getAppliedSchemaChangeEvents()` is empty 
> and the
> `MetadataApplier` is never invoked.
> h2. *Root Cause*
> 1. The MySQL enumerator DOES clean up its own state when the config no longer 
> matches `t`
>    (`MySqlSnapshotSplitAssigner#captureNewlyAddedTables`, removal branch), so 
> a re-added
>    `t` is correctly detected as newly-added, re-snapshotted, and the source 
> emits a
>    `CreateTableEvent` carrying the current (drifted) schema. This part works 
> as intended.
> 2. The `SchemaCoordinator`'s schema state, however, is never cleaned up:
>    `SchemaManager#originalSchemas` / `#evolvedSchemas` (checkpointed) have no 
> removal API
>    at all, so the stale V1 schema of the removed table survives across 
> savepoints.
> 3. `SchemaCoordinator#applySchemaChange` gates every event through
>    `SchemaUtils#isSchemaChangeEventRedundant`, whose `CreateTableEvent` 
> branch decides
>    "redundant" based only on the presence of *some* recorded schema:
>        createTableEvent -> {
>            // It has been applied if such table already exists
>            return latestSchema.isPresent();
>        },
>    The recorded schema is never compared with the incoming one, so the 
> `CreateTableEvent`
>    carrying the drifted schema is silently dropped:
>    - `MetadataApplier#applySchemaChange` is never invoked -> the downstream 
> table is not
>      created (symptom 1);
>    - neither the original nor the evolved schema state is updated, and the 
> stale evolved
>      schema is pushed back to `SchemaOperator` (`refreshedEvolvedSchemas`) 
> and to sink
>      writers (`DataSinkWriterOperator#emitLatestSchema`), so records get 
> coerced into the
>      stale shape and written against a table with a different schema -> 
> metadata mismatch
>      (symptom 2).
> The redundancy check was introduced by FLINK-36690 (restructured by 
> FLINK-37278) to skip
> duplicated `CreateTableEvent`s emitted in the snapshot stage. Skipping an 
> *identical*
> duplicate is correct; skipping one that carries a drifted schema against 
> stale state is not.
> h2. *Scope*
> - Regular schema evolution topology (default for MySQL and most sources); 
> 3.3.0-3.6.0 and
>   master are affected. 3.2 and earlier do not contain the redundancy check.
> - Any flow that re-snapshots a known table can hit this, including an ALTER 
> landing
>   mid-snapshot so that duplicated CreateTableEvents carry different schemas.
> - The distributed topology `SchemaCoordinator` has a similar gap 
> (`getLeastCommonSchema`
>   based derivation never narrows, and a `CreateTableEvent` is only emitted 
> when no evolved
>   schema exists); it can be tracked separately.
> - Known limitation remaining after the proposed fix: if the schema did NOT 
> drift (identical
>   `CreateTableEvent`) but the downstream table was dropped externally, the 
> framework cannot
>   distinguish it from a snapshot-stage duplicate and will not re-create it.
> h2. *Proposed Fix*
> 1. `SchemaUtils#isSchemaChangeEventRedundant`: treat a `CreateTableEvent` as 
> redundant only
>    when the recorded schema is identical to the incoming event's schema.
> 2. `SchemaCoordinator#deduceEvolvedSchemaChanges` (regular): when a 
> `CreateTableEvent`
>    targets a sink table that already has an evolved schema, emit
>    `CreateTableEvent(targetSchema)` (idempotent create-or-ignore, covering an 
> externally
>    dropped table) followed by alignment events from
>    `SchemaMergingUtils.getSchemaDifference(currentEvolvedSchema, 
> targetSchema)` instead of
>    forwarding the raw event. `targetSchema` honors the configured behavior: 
> the incoming
>    schema for EVOLVE / TRY_EVOLVE / EXCEPTION; `getLeastCommonSchema` (never 
> narrowing) for
>    LENIENT; the frozen current schema for IGNORE. Only the resulting 
> `CreateTableEvent` is
>    propagated downstream; alignment events are applied to the external system 
> only, and the
>    evolved-schema state update is guarded by the same redundancy check so 
> alignment events
>    are not applied to the state twice.
> 3. Tests: coordinator-level scenarios (EVOLVE with dropped/added column, 
> LENIENT, IGNORE,
>    identical-duplicate skip) plus unit tests for the utility change.
> h2. *Workaround (before a fixed release)*
> Create the downstream table manually with the schema recorded in the job 
> state - i.e. the
> OLD schema including the column(s) dropped upstream, as nullable - instead of 
> the current
> upstream schema. The pipeline coerces records into that schema (filling 
> dropped columns
> with null) and writes succeed.



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

Reply via email to