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]

Reply via email to