goutamadwant commented on code in PR #12284: URL: https://github.com/apache/seatunnel/pull/12284#discussion_r3997859744
########## 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: This ordering is not replay-safe when a rename or move uses a column added by the same composite event as its AFTER anchor. I reproduced this with input schema [a, b] and the valid sequence ADD c FIRST; CHANGE a TO x AFTER c. The translator emits the CHANGE before ADD c, and verifyReplay throws NoSuchElementException because c does not exist yet. A blanket ADD-before-CHANGE ordering is also insufficient because the inverse dependency can occur, for example CHANGE a TO b followed by ADD a. Please dependency-order or interleave additions and changes using their source names and positional anchors, and add the add-anchor case as a regression test. ########## 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: This fixture cannot exercise a FLOAT-to-DOUBLE change for the derived column. Existing ZetaSQLType numeric promotion produces DOUBLE for weight * 2 even when weight is FLOAT, so this assertion fails before the schema-change event is processed and double_weight remains DOUBLE after weight changes to DOUBLE. The corresponding E2E fixture has the same problem: it initially declares double_weight as FLOAT even though the transform produces DOUBLE, and the later source change cannot emit the expected derived-column type modification. Please use an expression and operand transition whose output type genuinely changes, such as an INT-to-BIGINT case, and align the initial sink schema with the transform pre-event output. ########## 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: This input event never sets sourceDialectName, so the translated event correctly carries null and the MySQL assertion at line 315 fails. I reproduced this test independently. Please set MySQL on the incoming modify event before passing it to the transform, as the composite-event helper does, so the test validates propagation of actual 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]
