DanielLeens commented on code in PR #12284: URL: https://github.com/apache/seatunnel/pull/12284#discussion_r3998097068
########## seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/sql/SQLSchemaChangeTranslator.java: ########## @@ -0,0 +1,604 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.transform.sql; + +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.Column; +import org.apache.seatunnel.api.table.catalog.ConstraintKey; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.api.table.schema.event.AlterColumnCommentEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableAddColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableChangeColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnsEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableDropColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableModifyColumnEvent; +import org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent; +import org.apache.seatunnel.api.table.schema.handler.AlterTableSchemaEventHandler; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.TreeMap; +import java.util.stream.Collectors; + +/** + * Translates upstream column-level DDL into events that describe the change of the SQL transform's + * own output. + * + * <p>Sinks advance the produced schema they received at planning time by applying the events they + * receive, so the emitted events must satisfy one contract: replaying them onto the pre-event + * produced schema with {@link AlterTableSchemaEventHandler} yields a schema equal to the produced + * schema after the event, and every output column keeps its data lineage. Attribution therefore + * works on input column identities (see {@link SQLLineageSchema}) and on output slots (see {@link + * SQLOutputSlot}), never on names alone. + * + * <p>All methods are static and stateless. Validation problems that a user can act on are reported + * as {@link IllegalArgumentException}; violations of internal invariants as {@link + * IllegalStateException}. + */ +final class SQLSchemaChangeTranslator { + + private SQLSchemaChangeTranslator() {} + + /** + * Column-level sub-events of a single or composite event; empty for table-level events. + * + * @param event any schema change event + * @return the ordered column-level events + */ + static List<AlterTableColumnEvent> flatten(SchemaChangeEvent event) { + if (event instanceof AlterTableColumnsEvent) { + return new ArrayList<>(((AlterTableColumnsEvent) event).getEvents()); + } + if (event instanceof AlterTableColumnEvent) { + return Collections.singletonList((AlterTableColumnEvent) event); + } + return Collections.emptyList(); + } + + /** + * Identities of the input columns that back the produced primary key, the produced constraint + * keys and the input partition keys. The shared handler never renames or removes a key column, + * so dropping or renaming one of these columns cannot be represented by any event. + * + * @param preSlots output slots before the event + * @param preOutput produced schema before the event + * @param partitionKeys input partition keys, may be null + * @param initial initial lineage of the input + * @return protected identities + */ + static Set<Integer> protectedIdentities( + List<SQLOutputSlot> preSlots, + TableSchema preOutput, + List<String> partitionKeys, + SQLLineageSchema initial) { + Set<String> protectedOutputNames = new HashSet<>(); + if (preOutput.getPrimaryKey() != null) { + protectedOutputNames.addAll(preOutput.getPrimaryKey().getColumnNames()); + } + if (preOutput.getConstraintKeys() != null) { + for (ConstraintKey constraintKey : preOutput.getConstraintKeys()) { + for (ConstraintKey.ConstraintKeyColumn keyColumn : constraintKey.getColumnNames()) { + protectedOutputNames.add(keyColumn.getColumnName()); + } + } + } + Set<Integer> identities = new HashSet<>(); + for (SQLOutputSlot slot : preSlots) { + if (!protectedOutputNames.contains(slot.getName())) { + continue; + } + for (String referenced : slot.getReferencedInputColumns()) { + int identity = initial.identityOf(referenced); + if (identity >= 0) { + identities.add(identity); + } + } + } + if (partitionKeys != null) { + for (String partitionKey : partitionKeys) { + int identity = initial.identityOf(partitionKey); + if (identity >= 0) { + identities.add(identity); + } + } + } + return identities; + } + + /** + * Rejects a drop or rename of a protected column. + * + * @param lineage lineage before the hint is applied + * @param hint the column-level event about to be applied + * @param protectedIdentities identities that must keep their name and existence + * @throws IllegalArgumentException when the hint drops or renames a protected column + */ + static void rejectProtectedColumnChange( + SQLLineageSchema lineage, + AlterTableColumnEvent hint, + Set<Integer> protectedIdentities) { + String target = null; + if (hint instanceof AlterTableDropColumnEvent) { + target = ((AlterTableDropColumnEvent) hint).getColumn(); + } else if (hint instanceof AlterTableChangeColumnEvent) { + AlterTableChangeColumnEvent change = (AlterTableChangeColumnEvent) hint; + if (!change.getOldColumn().equals(change.getColumn().getName())) { + target = change.getOldColumn(); + } + } + if (target != null && protectedIdentities.contains(lineage.identityOf(target))) { + throw new IllegalArgumentException( + String.format( + "column [%s] is part of the primary key, a constraint key or the partition keys and cannot be dropped or renamed", + target)); + } + } + + /** + * Computes the net effect of the change on every output slot and returns the events that + * reproduce it, ordered so that replaying them never collides on a column name. + * + * @param tableId produced table identifier used on the emitted events + * @param preSlots output slots before the event + * @param preOutput produced schema before the event + * @param preLineage initial lineage of the pre-event input + * @param finalSlots output slots after the event + * @param finalOutput produced schema after the event + * @param finalLineage lineage after all hints were applied + * @return ordered events, empty when the output did not change + * @throws IllegalArgumentException when renames form a cycle that no event order can replay + */ + static List<AlterTableColumnEvent> translate( + TableIdentifier tableId, + List<SQLOutputSlot> preSlots, + TableSchema preOutput, + SQLLineageSchema preLineage, + List<SQLOutputSlot> finalSlots, + TableSchema finalOutput, + SQLLineageSchema finalLineage) { + List<BoundSlot> pre = bind(preSlots, preOutput, preLineage); + List<BoundSlot> fin = bind(finalSlots, finalOutput, finalLineage); + + Map<Integer, BoundSlot> preStar = new LinkedHashMap<>(); + Map<Integer, BoundSlot> finalStar = new LinkedHashMap<>(); + Map<String, BoundSlot> preOther = new LinkedHashMap<>(); + Map<String, BoundSlot> finalOther = new LinkedHashMap<>(); + for (BoundSlot bound : pre) { + if (bound.isStar()) { + preStar.put(bound.identity(), bound); + } else { + preOther.put(bound.slot.pairingKey(), bound); + } + } + for (BoundSlot bound : fin) { + if (bound.isStar()) { + finalStar.put(bound.identity(), bound); + } else { + finalOther.put(bound.slot.pairingKey(), bound); + } + } + + List<AlterTableDropColumnEvent> drops = new ArrayList<>(); + List<AlterTableChangeColumnEvent> changes = new ArrayList<>(); + // Adds and modifies are keyed by final output index so AFTER anchors always exist. + TreeMap<Integer, AlterTableColumnEvent> adds = new TreeMap<>(); + TreeMap<Integer, AlterTableColumnEvent> modifies = new TreeMap<>(); + + for (BoundSlot preBound : preStar.values()) { + if (!finalStar.containsKey(preBound.identity())) { + drops.add(new AlterTableDropColumnEvent(tableId, preBound.column.getName())); + } + } + for (BoundSlot finalBound : finalStar.values()) { + BoundSlot preBound = preStar.get(finalBound.identity()); + if (preBound == null) { + Position position = + addPosition( + finalBound, + fin, + finalLineage.creatingHint(finalBound.identity()), + finalLineage); + adds.put( + finalBound.index, + new AlterTableAddColumnEvent( + tableId, finalBound.column, position.first, position.afterColumn)); + continue; + } + Position moved = movedPosition(preBound, finalBound, pre, fin, finalLineage); + if (!preBound.column.getName().equals(finalBound.column.getName())) { + changes.add( + new AlterTableChangeColumnEvent( + tableId, + preBound.column.getName(), + finalBound.column, + moved.first, + moved.afterColumn)); + } else if (!preBound.column.equals(finalBound.column) || moved.isSet()) { + modifies.put( + finalBound.index, + modifyOrComment(tableId, preBound.column, finalBound.column, moved)); + } + } + + for (Map.Entry<String, BoundSlot> entry : preOther.entrySet()) { + BoundSlot preBound = entry.getValue(); + BoundSlot finalBound = finalOther.get(entry.getKey()); + if (finalBound == null) { + throw new IllegalStateException( + String.format( + "output column [%s] disappeared from the query output", + preBound.column.getName())); + } + if (!preBound.signature.equals(finalBound.signature)) { + // The output column now carries data of a different physical input column. + drops.add(new AlterTableDropColumnEvent(tableId, preBound.column.getName())); + Position position = addPosition(finalBound, fin, null, finalLineage); + adds.put( + finalBound.index, + new AlterTableAddColumnEvent( + tableId, finalBound.column, position.first, position.afterColumn)); + } else if (!preBound.column.equals(finalBound.column)) { + modifies.put( + finalBound.index, + modifyOrComment( + tableId, preBound.column, finalBound.column, Position.none())); + } + } + for (String key : finalOther.keySet()) { + if (!preOther.containsKey(key)) { + throw new IllegalStateException( + String.format("output column [%s] appeared without a pre-event slot", key)); + } + } + + List<AlterTableColumnEvent> out = new ArrayList<>(drops); + out.addAll(orderChanges(changes, preOutput, drops)); Review Comment: Confirmed and fixed in 9bf19eb0efe. The translator no longer emits all renames before all adds. `translate` now drops first and then walks the final layout left to right on a simulated sink replay (`ReplayLayout`): the AFTER anchor of every positioned event is the final column to its left, which the walk has already placed under its final name, so `ADD c FIRST; CHANGE a TO x AFTER c` emits the add first and replays; a name that an add or a rename needs is freed just before by renaming its current holder in place (recursively for chains, a cycle still fails fast), so the inverse dependency `CHANGE a TO b; ADD a FIRST` replays as well. Regression tests: `testRenameAnchoredOnColumnAddedBySameEventReplays`, `testMoveAnchoredOnColumnAddedBySameEventReplays`, `testRenameFreesItsOldNameForAnAddAtALowerIndex` and `testBlockingRenameIsResolvedInPlaceFirst` in `SQLSchemaChangeTranslatorTest` (all go through `verifyReplay`), plus `testRenameAnchoredOnColumnAddedBySameCompositeReplays` throug h the transform in `SQLTransformSchemaChangeTest`. ########## seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/sql/SQLTransformSchemaChangeTest.java: ########## @@ -0,0 +1,891 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.transform.sql; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.event.EventType; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.PhysicalColumn; +import org.apache.seatunnel.api.table.catalog.PrimaryKey; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TablePath; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.api.table.schema.event.AlterTableAddColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableChangeColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnsEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableCommentEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableDropColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableModifyColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableNameEvent; +import org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent; +import org.apache.seatunnel.api.table.schema.handler.AlterTableSchemaEventHandler; +import org.apache.seatunnel.api.table.type.BasicType; +import org.apache.seatunnel.api.table.type.SeaTunnelDataType; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.transform.exception.TransformCommonErrorCode; +import org.apache.seatunnel.transform.exception.TransformException; +import org.apache.seatunnel.transform.rename.ConvertCase; +import org.apache.seatunnel.transform.rename.FieldRenameConfig; +import org.apache.seatunnel.transform.rename.FieldRenameTransform; +import org.apache.seatunnel.transform.sql.zeta.ZetaSQLEngine; +import org.apache.seatunnel.transform.sql.zeta.ZetaUDF; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Covers the SQL transform's schema change translation end to end: output-relative events for star + * projections, projections, aliases and derived columns, absorbed changes, fail-fast cases, lineage + * through composites, staged engine hand-offs at chain positions greater than zero, the multi-table + * wrapper with several rules per table, resynchronisation events and the engine and UDF lifecycle + * across changes. + */ +public class SQLTransformSchemaChangeTest { + + private static final TablePath TBL = TablePath.of("db", "products"); + private static final TableIdentifier TID = TableIdentifier.of("catalog", TBL); + + @Test + public void testStarAddColumnKeepsUpstreamShape() { + SQLTransform transform = transform("select * from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + AlterTableColumnsEvent event = composite(addAge()); + + SchemaChangeEvent out = transform.mapSchemaChangeEvent(event); + + Assertions.assertTrue(out instanceof AlterTableColumnsEvent); + List<AlterTableColumnEvent> events = ((AlterTableColumnsEvent) out).getEvents(); + Assertions.assertEquals(1, events.size()); + AlterTableAddColumnEvent add = (AlterTableAddColumnEvent) events.get(0); + Assertions.assertEquals("age", add.getColumn().getName()); + Assertions.assertFalse(add.isFirst()); + Assertions.assertNull(add.getAfterColumn()); + Assertions.assertEquals("MySQL", add.getSourceDialectName()); + Assertions.assertEquals("job-1", add.getJobId()); + Assertions.assertEquals(event.getStatement(), add.getStatement()); + Assertions.assertSame(transform.getProducedCatalogTable(), add.getChangeAfter()); + Assertions.assertSame(transform.getProducedCatalogTable(), out.getChangeAfter()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d, 20})); + Assertions.assertEquals(4, rows.get(0).getArity()); + Assertions.assertEquals(20, rows.get(0).getField(3)); + } + + @Test + public void testStarAddAfterAndFirstKeepPositions() { + SQLTransform after = transform("select * from products", baseTable()); + AlterTableAddColumnEvent addAfter = + (AlterTableAddColumnEvent) + after.mapSchemaChangeEvent( + AlterTableAddColumnEvent.addAfter( + TID, column("age", BasicType.INT_TYPE, "int"), "id")); + Assertions.assertEquals("id", addAfter.getAfterColumn()); + Assertions.assertArrayEquals( + new String[] {"id", "age", "name", "weight"}, + after.getProducedCatalogTable().getTableSchema().getFieldNames()); + + SQLTransform first = transform("select * from products", baseTable()); + AlterTableAddColumnEvent addFirst = + (AlterTableAddColumnEvent) + first.mapSchemaChangeEvent( + AlterTableAddColumnEvent.addFirst( + TID, column("age", BasicType.INT_TYPE, "int"))); + Assertions.assertTrue(addFirst.isFirst()); + } + + @Test + public void testStarWithTrailingExpressionAddsAfterLastStarColumn() { + SQLTransform transform = + transform("select *, weight * 2 as double_weight from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + SchemaChangeEvent out = transform.mapSchemaChangeEvent(composite(addAge())); + + AlterTableAddColumnEvent add = + (AlterTableAddColumnEvent) ((AlterTableColumnsEvent) out).getEvents().get(0); + Assertions.assertEquals("weight", add.getAfterColumn()); + Assertions.assertArrayEquals( + new String[] {"id", "name", "weight", "age", "double_weight"}, + transform.getProducedCatalogTable().getTableSchema().getFieldNames()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.5d, 20})); + Assertions.assertEquals(5, rows.get(0).getArity()); + Assertions.assertEquals(20, rows.get(0).getField(3)); + Assertions.assertEquals(3.0d, ((Number) rows.get(0).getField(4)).doubleValue(), 0.0001d); + } + + @Test + public void testProjectionAbsorbsUnrelatedAddAndDrop() { + SQLTransform transform = transform("select id, name from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + Assertions.assertNull(transform.mapSchemaChangeEvent(composite(addAge()))); + Assertions.assertEquals(before, transform.getProducedCatalogTable().getTableSchema()); + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d, 20})); + Assertions.assertEquals(2, rows.get(0).getArity()); + Assertions.assertEquals("a", rows.get(0).getField(1)); + + Assertions.assertNull( + transform.mapSchemaChangeEvent( + composite(new AlterTableDropColumnEvent(TID, "age")))); + rows = transform.flatMap(new SeaTunnelRow(new Object[] {2L, "b", 1.0d})); + Assertions.assertEquals(2, rows.get(0).getArity()); + } + + @Test + public void testDropOrRenameOfReferencedColumnFailsFastAndLeavesStateUntouched() { + SQLTransform transform = transform("select id, name from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + TransformException drop = + Assertions.assertThrows( + TransformException.class, + () -> + transform.mapSchemaChangeEvent( + composite(new AlterTableDropColumnEvent(TID, "name")))); + Assertions.assertEquals( + TransformCommonErrorCode.SQL_SCHEMA_CHANGE_INCOMPATIBLE, + drop.getSeaTunnelErrorCode()); + Assertions.assertTrue(drop.getMessage().contains("[name]"), drop.getMessage()); + + TransformException rename = + Assertions.assertThrows( + TransformException.class, + () -> + transform.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column( + "full_name", + BasicType.STRING_TYPE, + "varchar(255)"))))); + Assertions.assertEquals( + TransformCommonErrorCode.SQL_SCHEMA_CHANGE_INCOMPATIBLE, + rename.getSeaTunnelErrorCode()); + + Assertions.assertEquals(before, transform.getProducedCatalogTable().getTableSchema()); + List<SeaTunnelRow> rows = transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d})); + Assertions.assertEquals(2, rows.get(0).getArity()); + } + + @Test + public void testStarRenameEmitsChangeWithFinalColumn() { + SQLTransform transform = transform("select * from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + PhysicalColumn fullName = column("full_name", BasicType.STRING_TYPE, "varchar(255)"); + + SchemaChangeEvent out = + transform.mapSchemaChangeEvent( + composite(AlterTableChangeColumnEvent.change(TID, "name", fullName))); + + AlterTableChangeColumnEvent change = + (AlterTableChangeColumnEvent) ((AlterTableColumnsEvent) out).getEvents().get(0); + Assertions.assertEquals("name", change.getOldColumn()); + Assertions.assertEquals("full_name", change.getColumn().getName()); + Assertions.assertEquals("varchar(255)", change.getColumn().getSourceType()); + Assertions.assertEquals("MySQL", change.getSourceDialectName()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + } + + @Test + public void testCompositeRenameReuseAndDropReaddKeepLineage() { + SQLTransform star = transform("select * from products", baseTable()); + TableSchema starBefore = star.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent starOut = + star.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("name_old", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> starEvents = ((AlterTableColumnsEvent) starOut).getEvents(); + Assertions.assertEquals(2, starEvents.size()); + Assertions.assertTrue(starEvents.get(0) instanceof AlterTableChangeColumnEvent); + Assertions.assertTrue(starEvents.get(1) instanceof AlterTableAddColumnEvent); + assertReplays(starBefore, starOut, star.getProducedCatalogTable().getTableSchema()); + + SQLTransform reference = transform("select id, name from products", baseTable()); + TableSchema referenceBefore = reference.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent referenceOut = + reference.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("name_old", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> referenceEvents = + ((AlterTableColumnsEvent) referenceOut).getEvents(); + Assertions.assertEquals(2, referenceEvents.size()); + Assertions.assertTrue(referenceEvents.get(0) instanceof AlterTableDropColumnEvent); + Assertions.assertEquals( + "name", ((AlterTableDropColumnEvent) referenceEvents.get(0)).getColumn()); + AlterTableAddColumnEvent readd = (AlterTableAddColumnEvent) referenceEvents.get(1); + Assertions.assertEquals("name", readd.getColumn().getName()); + Assertions.assertEquals("varchar(64)", readd.getColumn().getSourceType()); + assertReplays( + referenceBefore, + referenceOut, + reference.getProducedCatalogTable().getTableSchema()); + + SQLTransform dropReadd = transform("select id, name from products", baseTable()); + TableSchema dropReaddBefore = dropReadd.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent dropReaddOut = + dropReadd.mapSchemaChangeEvent( + composite( + new AlterTableDropColumnEvent(TID, "name"), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> dropReaddEvents = + ((AlterTableColumnsEvent) dropReaddOut).getEvents(); + Assertions.assertTrue(dropReaddEvents.get(0) instanceof AlterTableDropColumnEvent); + Assertions.assertTrue(dropReaddEvents.get(1) instanceof AlterTableAddColumnEvent); + assertReplays( + dropReaddBefore, + dropReaddOut, + dropReadd.getProducedCatalogTable().getTableSchema()); + + SQLTransform roundTrip = transform("select id, name from products", baseTable()); + Assertions.assertNull( + roundTrip.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("tmp", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableChangeColumnEvent.change( + TID, + "tmp", + column("name", BasicType.STRING_TYPE, "varchar(255)"))))); + } + + @Test + public void testModifyDirectReferenceCarriesDialectAndSourceType() { + SQLTransform transform = transform("select id, name as n from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + SchemaChangeEvent out = + transform.mapSchemaChangeEvent( + AlterTableModifyColumnEvent.modify( + TID, column("name", BasicType.STRING_TYPE, "longtext"))); + + Assertions.assertTrue(out instanceof AlterTableModifyColumnEvent); + AlterTableModifyColumnEvent modify = (AlterTableModifyColumnEvent) out; + Assertions.assertEquals("n", modify.getColumn().getName()); + Assertions.assertEquals("longtext", modify.getColumn().getSourceType()); + Assertions.assertEquals("MySQL", modify.getSourceDialectName()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + } + + @Test + public void testModifyDerivedColumnHasNoDialectAndNoSourceType() { + TableSchema schema = + TableSchema.builder() + .column(column("id", BasicType.LONG_TYPE, "bigint")) + .column(column("weight", BasicType.FLOAT_TYPE, "float")) + .primaryKey(PrimaryKey.of("pk", Collections.singletonList("id"))) + .build(); + SQLTransform transform = + transform( + "select id, weight, weight * 2 as double_weight from products", + table(schema)); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + Assertions.assertEquals( Review Comment: Confirmed and fixed in 9bf19eb0efe. The unit test now uses `quantity INT` with `quantity * 2 as double_quantity` (derived INT) and a `MODIFY quantity BIGINT`, after which the expression is derived as BIGINT, so the derived-column modify is real; it also asserts the post-change row computes 42 for 21. The E2E fixture keeps the `weight * 2` projection but the sink now starts with `double_weight DOUBLE` (the type the transform derives for the FLOAT operand) and `modify_weight_type.sql` changes `weight` to `DECIMAL(12,3)`, which turns the derived type from DOUBLE into DECIMAL(12,3); the test asserts `decimal(12,3)` for both `weight` and `double_weight`. The docs example sentence was corrected the same way, and the fixture comment explains why a change to DOUBLE would not retype the expression. ########## seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/sql/SQLTransformSchemaChangeTest.java: ########## @@ -0,0 +1,891 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.transform.sql; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.event.EventType; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.PhysicalColumn; +import org.apache.seatunnel.api.table.catalog.PrimaryKey; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TablePath; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.api.table.schema.event.AlterTableAddColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableChangeColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableColumnsEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableCommentEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableDropColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableModifyColumnEvent; +import org.apache.seatunnel.api.table.schema.event.AlterTableNameEvent; +import org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent; +import org.apache.seatunnel.api.table.schema.handler.AlterTableSchemaEventHandler; +import org.apache.seatunnel.api.table.type.BasicType; +import org.apache.seatunnel.api.table.type.SeaTunnelDataType; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.transform.exception.TransformCommonErrorCode; +import org.apache.seatunnel.transform.exception.TransformException; +import org.apache.seatunnel.transform.rename.ConvertCase; +import org.apache.seatunnel.transform.rename.FieldRenameConfig; +import org.apache.seatunnel.transform.rename.FieldRenameTransform; +import org.apache.seatunnel.transform.sql.zeta.ZetaSQLEngine; +import org.apache.seatunnel.transform.sql.zeta.ZetaUDF; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Covers the SQL transform's schema change translation end to end: output-relative events for star + * projections, projections, aliases and derived columns, absorbed changes, fail-fast cases, lineage + * through composites, staged engine hand-offs at chain positions greater than zero, the multi-table + * wrapper with several rules per table, resynchronisation events and the engine and UDF lifecycle + * across changes. + */ +public class SQLTransformSchemaChangeTest { + + private static final TablePath TBL = TablePath.of("db", "products"); + private static final TableIdentifier TID = TableIdentifier.of("catalog", TBL); + + @Test + public void testStarAddColumnKeepsUpstreamShape() { + SQLTransform transform = transform("select * from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + AlterTableColumnsEvent event = composite(addAge()); + + SchemaChangeEvent out = transform.mapSchemaChangeEvent(event); + + Assertions.assertTrue(out instanceof AlterTableColumnsEvent); + List<AlterTableColumnEvent> events = ((AlterTableColumnsEvent) out).getEvents(); + Assertions.assertEquals(1, events.size()); + AlterTableAddColumnEvent add = (AlterTableAddColumnEvent) events.get(0); + Assertions.assertEquals("age", add.getColumn().getName()); + Assertions.assertFalse(add.isFirst()); + Assertions.assertNull(add.getAfterColumn()); + Assertions.assertEquals("MySQL", add.getSourceDialectName()); + Assertions.assertEquals("job-1", add.getJobId()); + Assertions.assertEquals(event.getStatement(), add.getStatement()); + Assertions.assertSame(transform.getProducedCatalogTable(), add.getChangeAfter()); + Assertions.assertSame(transform.getProducedCatalogTable(), out.getChangeAfter()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d, 20})); + Assertions.assertEquals(4, rows.get(0).getArity()); + Assertions.assertEquals(20, rows.get(0).getField(3)); + } + + @Test + public void testStarAddAfterAndFirstKeepPositions() { + SQLTransform after = transform("select * from products", baseTable()); + AlterTableAddColumnEvent addAfter = + (AlterTableAddColumnEvent) + after.mapSchemaChangeEvent( + AlterTableAddColumnEvent.addAfter( + TID, column("age", BasicType.INT_TYPE, "int"), "id")); + Assertions.assertEquals("id", addAfter.getAfterColumn()); + Assertions.assertArrayEquals( + new String[] {"id", "age", "name", "weight"}, + after.getProducedCatalogTable().getTableSchema().getFieldNames()); + + SQLTransform first = transform("select * from products", baseTable()); + AlterTableAddColumnEvent addFirst = + (AlterTableAddColumnEvent) + first.mapSchemaChangeEvent( + AlterTableAddColumnEvent.addFirst( + TID, column("age", BasicType.INT_TYPE, "int"))); + Assertions.assertTrue(addFirst.isFirst()); + } + + @Test + public void testStarWithTrailingExpressionAddsAfterLastStarColumn() { + SQLTransform transform = + transform("select *, weight * 2 as double_weight from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + SchemaChangeEvent out = transform.mapSchemaChangeEvent(composite(addAge())); + + AlterTableAddColumnEvent add = + (AlterTableAddColumnEvent) ((AlterTableColumnsEvent) out).getEvents().get(0); + Assertions.assertEquals("weight", add.getAfterColumn()); + Assertions.assertArrayEquals( + new String[] {"id", "name", "weight", "age", "double_weight"}, + transform.getProducedCatalogTable().getTableSchema().getFieldNames()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.5d, 20})); + Assertions.assertEquals(5, rows.get(0).getArity()); + Assertions.assertEquals(20, rows.get(0).getField(3)); + Assertions.assertEquals(3.0d, ((Number) rows.get(0).getField(4)).doubleValue(), 0.0001d); + } + + @Test + public void testProjectionAbsorbsUnrelatedAddAndDrop() { + SQLTransform transform = transform("select id, name from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + Assertions.assertNull(transform.mapSchemaChangeEvent(composite(addAge()))); + Assertions.assertEquals(before, transform.getProducedCatalogTable().getTableSchema()); + List<SeaTunnelRow> rows = + transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d, 20})); + Assertions.assertEquals(2, rows.get(0).getArity()); + Assertions.assertEquals("a", rows.get(0).getField(1)); + + Assertions.assertNull( + transform.mapSchemaChangeEvent( + composite(new AlterTableDropColumnEvent(TID, "age")))); + rows = transform.flatMap(new SeaTunnelRow(new Object[] {2L, "b", 1.0d})); + Assertions.assertEquals(2, rows.get(0).getArity()); + } + + @Test + public void testDropOrRenameOfReferencedColumnFailsFastAndLeavesStateUntouched() { + SQLTransform transform = transform("select id, name from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + TransformException drop = + Assertions.assertThrows( + TransformException.class, + () -> + transform.mapSchemaChangeEvent( + composite(new AlterTableDropColumnEvent(TID, "name")))); + Assertions.assertEquals( + TransformCommonErrorCode.SQL_SCHEMA_CHANGE_INCOMPATIBLE, + drop.getSeaTunnelErrorCode()); + Assertions.assertTrue(drop.getMessage().contains("[name]"), drop.getMessage()); + + TransformException rename = + Assertions.assertThrows( + TransformException.class, + () -> + transform.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column( + "full_name", + BasicType.STRING_TYPE, + "varchar(255)"))))); + Assertions.assertEquals( + TransformCommonErrorCode.SQL_SCHEMA_CHANGE_INCOMPATIBLE, + rename.getSeaTunnelErrorCode()); + + Assertions.assertEquals(before, transform.getProducedCatalogTable().getTableSchema()); + List<SeaTunnelRow> rows = transform.flatMap(new SeaTunnelRow(new Object[] {1L, "a", 1.0d})); + Assertions.assertEquals(2, rows.get(0).getArity()); + } + + @Test + public void testStarRenameEmitsChangeWithFinalColumn() { + SQLTransform transform = transform("select * from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + PhysicalColumn fullName = column("full_name", BasicType.STRING_TYPE, "varchar(255)"); + + SchemaChangeEvent out = + transform.mapSchemaChangeEvent( + composite(AlterTableChangeColumnEvent.change(TID, "name", fullName))); + + AlterTableChangeColumnEvent change = + (AlterTableChangeColumnEvent) ((AlterTableColumnsEvent) out).getEvents().get(0); + Assertions.assertEquals("name", change.getOldColumn()); + Assertions.assertEquals("full_name", change.getColumn().getName()); + Assertions.assertEquals("varchar(255)", change.getColumn().getSourceType()); + Assertions.assertEquals("MySQL", change.getSourceDialectName()); + assertReplays(before, out, transform.getProducedCatalogTable().getTableSchema()); + } + + @Test + public void testCompositeRenameReuseAndDropReaddKeepLineage() { + SQLTransform star = transform("select * from products", baseTable()); + TableSchema starBefore = star.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent starOut = + star.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("name_old", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> starEvents = ((AlterTableColumnsEvent) starOut).getEvents(); + Assertions.assertEquals(2, starEvents.size()); + Assertions.assertTrue(starEvents.get(0) instanceof AlterTableChangeColumnEvent); + Assertions.assertTrue(starEvents.get(1) instanceof AlterTableAddColumnEvent); + assertReplays(starBefore, starOut, star.getProducedCatalogTable().getTableSchema()); + + SQLTransform reference = transform("select id, name from products", baseTable()); + TableSchema referenceBefore = reference.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent referenceOut = + reference.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("name_old", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> referenceEvents = + ((AlterTableColumnsEvent) referenceOut).getEvents(); + Assertions.assertEquals(2, referenceEvents.size()); + Assertions.assertTrue(referenceEvents.get(0) instanceof AlterTableDropColumnEvent); + Assertions.assertEquals( + "name", ((AlterTableDropColumnEvent) referenceEvents.get(0)).getColumn()); + AlterTableAddColumnEvent readd = (AlterTableAddColumnEvent) referenceEvents.get(1); + Assertions.assertEquals("name", readd.getColumn().getName()); + Assertions.assertEquals("varchar(64)", readd.getColumn().getSourceType()); + assertReplays( + referenceBefore, + referenceOut, + reference.getProducedCatalogTable().getTableSchema()); + + SQLTransform dropReadd = transform("select id, name from products", baseTable()); + TableSchema dropReaddBefore = dropReadd.getProducedCatalogTable().getTableSchema(); + SchemaChangeEvent dropReaddOut = + dropReadd.mapSchemaChangeEvent( + composite( + new AlterTableDropColumnEvent(TID, "name"), + AlterTableAddColumnEvent.add( + TID, + column("name", BasicType.STRING_TYPE, "varchar(64)")))); + List<AlterTableColumnEvent> dropReaddEvents = + ((AlterTableColumnsEvent) dropReaddOut).getEvents(); + Assertions.assertTrue(dropReaddEvents.get(0) instanceof AlterTableDropColumnEvent); + Assertions.assertTrue(dropReaddEvents.get(1) instanceof AlterTableAddColumnEvent); + assertReplays( + dropReaddBefore, + dropReaddOut, + dropReadd.getProducedCatalogTable().getTableSchema()); + + SQLTransform roundTrip = transform("select id, name from products", baseTable()); + Assertions.assertNull( + roundTrip.mapSchemaChangeEvent( + composite( + AlterTableChangeColumnEvent.change( + TID, + "name", + column("tmp", BasicType.STRING_TYPE, "varchar(255)")), + AlterTableChangeColumnEvent.change( + TID, + "tmp", + column("name", BasicType.STRING_TYPE, "varchar(255)"))))); + } + + @Test + public void testModifyDirectReferenceCarriesDialectAndSourceType() { + SQLTransform transform = transform("select id, name as n from products", baseTable()); + TableSchema before = transform.getProducedCatalogTable().getTableSchema(); + + SchemaChangeEvent out = + transform.mapSchemaChangeEvent( + AlterTableModifyColumnEvent.modify( Review Comment: Confirmed and fixed in 9bf19eb0efe. Added a `single(event)` helper that sets the job id, statement and `MySQL` dialect the way CDC sources do, mirroring the composite helper, and applied it to every single incoming event in `SQLTransformSchemaChangeTest`, so the assertion validates propagation of real input metadata. -- 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]
