DanielLeens opened a new pull request, #12284:
URL: https://github.com/apache/seatunnel/pull/12284
### Purpose of this pull request
Closes #8253.
Today the `Sql` transform reacts to a `SchemaChangeEvent` by rebuilding its
input schema and discarding its cached engine, but it forwards the **upstream**
event downstream unchanged (`SQLTransform.mapSchemaChangeEvent`). Sinks
re-derive their schema and the physical DDL from the column sub-events of that
event, so as soon as the query is anything other than a bare `select *` the
sink diverges from the rows the transform emits:
- `select id, name` plus an upstream `ADD COLUMN age`: the sink table gains
`age` and the sink writer expects 3 fields while the transform keeps emitting
2-field rows.
- `select id, cast(price as decimal(10,2)) as price` plus `MODIFY price
DOUBLE`: the sink column becomes `DOUBLE` while the rows still carry decimals.
- `select id, name` plus `DROP COLUMN name`: nothing fails at the event;
every following row fails in `ZetaSQLFunction` and is classified as a row
error, so with row error skipping the data is silently lost.
This PR makes the SQL transform translate every column-level change into a
change of its **own output**:
- **Identity-tracked lineage.** Sub-events are applied in order to an
identity-annotated copy of the input (`SQLLineageSchema`); output columns are
bound to input identities (`SQLOutputSlot`); the emitted events are the net
effect per output column (`SQLSchemaChangeTranslator`). `CHANGE a -> b; ADD a`
becomes a rename plus an add, `DROP a; ADD a` becomes a drop plus an add (never
a modify), a rename reverted in the same statement is absorbed, and `select a`
with `DROP a; ADD a BIGINT` becomes a drop plus an add of the output column.
- **Replay contract.** Replaying the emitted events with
`AlterTableSchemaEventHandler` onto the pre-event produced schema yields a
`TableSchema` that is `equals` to the produced schema after the event (all
column attributes, primary key, constraint keys). This is checked at runtime
and in every unit test.
- **Transactional evaluation.** The final input is evaluated on an isolated
engine; live state (`inputCatalogTable`, engine, row type, produced table) is
swapped only after every check passed, and the retired engine is closed exactly
once, so UDF resources no longer leak on schema changes. A rejected event
leaves the transform untouched at any chain position.
- **Staged engine hand-off.** `setInputCatalogTable` only stages the
pre-event state and the handed-over input; `mapSchemaChangeEvent` validates and
commits. This also covers `rule_match_mode = ALL_MATCH` chains where the same
transform is handed its input twice for one event.
- **Provenance-aware metadata.** `sourceDialectName` is copied to an
`ADD`/`MODIFY`/`CHANGE` sub-event only when the output column carries a
`sourceType` inherited from an input column, so JDBC dialects never emit a
`null` same-dialect type for derived expression columns and reconvert their
SeaTunnel type instead.
- **Fail fast, scoped.** A dropped or renamed column that the query
references, a drop or rename of a primary key, constraint key or partition
column, a rename cycle, a duplicate output name, an ordering comparison in
`WHERE` whose operand families no longer match, or an upstream produced schema
that does not match the event fail the job at the event with the new
`TRANSFORM_COMMON-09` error code. Type changes that only surface inside
functions, UDFs or lateral view arguments still surface at row time, as
documented.
- **Identity path of multi-table wrappers.** `IdentityMapTransform` and
`IdentityFlatMapTransform` now derive their produced table from their live
input and adopt the event's `changeAfter`, so a multi-table SQL job no longer
stamps a stale table into the event for tables that no rule matches.
`select *` jobs keep emitting the same DDL as today (same names, same
positions, same metadata).
### Does this PR introduce _any_ user-facing change?
- The `Sql` transform now supports schema evolution. Behaviour is documented
in `docs/en/transforms/sql.md` and `docs/zh/transforms/sql.md` (rules, lineage,
limits).
- New error code `TRANSFORM_COMMON-09 SQL_SCHEMA_CHANGE_INCOMPATIBLE`.
- Jobs with a non-`select *` SQL transform and `schema-changes.enabled =
true` were broken at the sink before; they now either receive correct events or
fail fast with an actionable message. Chains where a transform placed before
the SQL transform forwards an event that does not describe its own produced
schema (for example `Metadata` or `Copy` followed by `Sql`) now fail with
`TRANSFORM_COMMON-09` instead of writing misaligned rows.
- `docs/{en,zh}/introduction/configuration/schema-evolution.md` replace the
"schema evolution does not support transforms" note with a per-transform
support matrix derived from the transform class inventory.
- No configuration option is added or changed. `seatunnel-api`, the engine
and the sink side are untouched.
### How was this PR tested?
Unit tests (`seatunnel-transforms-v2`):
- `SQLTransformSchemaChangeTest`: star projections, projections, aliases,
derived columns, absorbed changes, fail-fast cases with untouched state,
composite rename-and-reuse / drop-and-re-add / rename-away-and-back, key column
protection, duplicate output names, `WHERE` type compatibility, table-level
events, resynchronisation from an event that carries the whole table, staged
hand-offs at chain positions greater than zero (translation equal to position
zero, rejected events, upstream mismatch, rows during a pending hand-off),
`rule_match_mode = ALL_MATCH` at engine position greater than zero, a
`FieldRename -> Sql` chain, and engine plus UDF lifecycle counts across
successful, rejected and staged changes.
- `SQLSchemaChangeTranslatorTest`: lineage, ordering, position derivation,
key protection, replay verification, event shape and metadata rules.
- `ZetaSQLEngineReferencedColumnsTest`: referenced column analysis, output
slot description, `WHERE` type checks.
- `SQLMultiCatalogSchemaChangeTest` and
`MetadataMultiCatalogSchemaChangeTest`: unmatched-table identity path including
a rename handed over before the event.
- `TransformChainLiveAlterTest`: now asserts the outgoing SQL event (`ADD
... AFTER <last star column>`).
E2E (`connector-cdc-mysql-e2e`, `MysqlCDCWithSchemaChangeIT`):
- `mysqlcdc_to_mysql_with_schema_change_sql_star.conf`: the existing
add/drop/change/modify/comment sequence through a `select *` transform, plus a
new `rename_reuse_columns` template (rename with name reuse, then drop and
re-create, each in one statement).
- `mysqlcdc_to_mysql_with_schema_change_sql_projection.conf`: `select id,
name, weight, weight * 2 as double_weight`; unrelated add/drop/rename absorbed,
`modify name longtext` forwarded with the source type, a new
`modify_weight_type` template changing `weight` to `DOUBLE` (direct column and
derived expression column both become `DOUBLE` in the same-dialect MySQL sink),
and a new `drop_readd_projected` template dropping and re-creating the
referenced `name` column in one statement.
Validation policy: only `./mvnw spotless:apply` was run locally on the
touched modules; compilation, unit tests and E2E are validated by this PR's
GitHub CI.
### Dependencies and follow-ups
- Recovery after a failover depends on #11503: transforms are rebuilt from
the planning-time schema on restore and the source's restore event is what
resynchronises them. This PR handles any `AlterTableEvent` that carries the
whole table in `changeAfter` as a resynchronisation (no DDL derived), which is
exactly what `RestoreTableSchemaEvent` of #11503 is, so the two changes compose
regardless of merge order. The "DDL, savepoint, restore, then rows with no
further DDL" E2E is added once #11503 merges. Without #11503 a job restored
after an output-affecting DDL now fails fast at the next DDL with
`TRANSFORM_COMMON-09` instead of silently writing against a stale schema.
- Follow-up: an explicit "input already handed over for this event" contract
for `AbstractCatalogSupportMapTransform`, so map transforms placed after
another transform no longer re-apply a column rename (pre-existing
`IndexOutOfBoundsException` in
`AlterTableSchemaEventHandler.applyChangeColumn`).
- Follow-up: `Calcite`, `Python` and `TextChunk` can adopt the translator.
### Check list
* [x] Code changed are covered with tests, or it does not need tests for
reason: covered by unit tests and E2E listed above
* [x] If any new Jar binary package adding in your PR, please add License
Notice according [New License
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/contribution/new-license.md):
no new dependency
* [x] If necessary, please update the documentation to describe the new
feature. https://github.com/apache/seatunnel/tree/dev/docs
* [x] If you are contributing the connector code, please check that the
following files are updated: not a connector change
* [ ] Update the
[`release-note`](https://github.com/apache/seatunnel/blob/dev/release-note.md):
the file does not exist on `dev`, so no entry was added.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]