This is an automated email from the ASF dual-hosted git repository.

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new d6df7f2f25 [Fix][Connector-V2] Exclude dropped columns in postgresql 
connector (#11358)
d6df7f2f25 is described below

commit d6df7f2f25abface0aa9498a8fc7b3a18572c32e
Author: NganWave <[email protected]>
AuthorDate: Wed Jul 15 12:02:04 2026 +0800

    [Fix][Connector-V2] Exclude dropped columns in postgresql connector (#11358)
    
    Co-authored-by: DanielLeens <[email protected]>
---
 .../jdbc/catalog/psql/PostgresCatalog.java         |  3 ++
 .../connectors/seatunnel/jdbc/JdbcPostgresIT.java  | 57 ++++++++++++++++++++++
 .../connectors/seatunnel/jdbc/JdbcOpenGaussIT.java | 43 ++++++++++++++++
 3 files changed, 103 insertions(+)

diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
index ccb895edc4..fd07f04029 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
@@ -78,6 +78,9 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
                     + "    n.nspname = '%s'\n"
                     + "    AND c.relname = '%s'\n"
                     + "    AND a.attnum > 0\n"
+                    // PostgreSQL-compatible catalogs retain placeholder 
attributes after DROP
+                    // COLUMN.
+                    + "    AND NOT a.attisdropped\n"
                     + "ORDER BY \n"
                     + "    a.attnum;";
 
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcPostgresIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcPostgresIT.java
index bc469dca41..0d1978cf9a 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcPostgresIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcPostgresIT.java
@@ -22,6 +22,7 @@ import 
org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
 
 import org.apache.seatunnel.api.table.catalog.Catalog;
 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.PrimaryKey;
 import org.apache.seatunnel.api.table.catalog.TablePath;
@@ -55,8 +56,10 @@ import java.sql.Connection;
 import java.sql.DriverManager;
 import java.sql.SQLException;
 import java.sql.Statement;
+import java.util.Arrays;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
 import static org.awaitility.Awaitility.given;
@@ -70,6 +73,13 @@ public class JdbcPostgresIT extends TestSuiteBase implements 
TestResource {
             
"https://repo1.maven.org/maven2/net/postgis/postgis-jdbc/2.5.1/postgis-jdbc-2.5.1.jar";;
     private static final String PG_GEOMETRY_JAR =
             
"https://repo1.maven.org/maven2/net/postgis/postgis-geometry/2.5.1/postgis-geometry-2.5.1.jar";;
+
+    /**
+     * PostgreSQL table used to verify catalog behavior after dropping a 
physical column while
+     * preserving the order of all remaining visible columns.
+     */
+    private static final String DROPPED_COLUMN_TABLE = 
"pg_catalog_dropped_column_test";
+
     private static final List<String> PG_CONFIG_FILE_LIST =
             Lists.newArrayList(
                     "/jdbc_postgres_source_and_sink.conf",
@@ -397,6 +407,53 @@ public class JdbcPostgresIT extends TestSuiteBase 
implements TestResource {
         catalog.close();
     }
 
+    /**
+     * Verifies PostgreSQL catalog discovery preserves visible column order 
after a column is
+     * dropped.
+     */
+    @Test
+    public void testCatalogExcludesDroppedColumns() {
+        String schema = "public";
+        String databaseName = POSTGRESQL_CONTAINER.getDatabaseName();
+        TablePath tablePath = TablePath.of(databaseName, schema, 
DROPPED_COLUMN_TABLE);
+        PostgresCatalog postgresCatalog =
+                new PostgresCatalog(
+                        DatabaseIdentifier.POSTGRESQL,
+                        POSTGRESQL_CONTAINER.getUsername(),
+                        POSTGRESQL_CONTAINER.getPassword(),
+                        
JdbcUrlUtil.getUrlInfo(POSTGRESQL_CONTAINER.getJdbcUrl()),
+                        schema,
+                        null);
+        postgresCatalog.open();
+        try {
+            postgresCatalog.dropTable(tablePath, true);
+            postgresCatalog.executeSql(
+                    tablePath,
+                    "CREATE TABLE "
+                            + schema
+                            + "."
+                            + DROPPED_COLUMN_TABLE
+                            + " (c1 INT, c2 VARCHAR(50), c3 INT, c4 TEXT)");
+            postgresCatalog.executeSql(
+                    tablePath,
+                    "ALTER TABLE " + schema + "." + DROPPED_COLUMN_TABLE + " 
DROP COLUMN c2");
+
+            CatalogTable table = postgresCatalog.getTable(tablePath);
+            List<String> actualColumns =
+                    table.getTableSchema().getColumns().stream()
+                            .map(Column::getName)
+                            .collect(Collectors.toList());
+
+            Assertions.assertEquals(Arrays.asList("c1", "c3", "c4"), 
actualColumns);
+        } finally {
+            try {
+                postgresCatalog.dropTable(tablePath, true);
+            } finally {
+                postgresCatalog.close();
+            }
+        }
+    }
+
     private void initializeJdbcTable() {
         try (Connection connection = getJdbcConnection()) {
             Statement statement = connection.createStatement();
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcOpenGaussIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcOpenGaussIT.java
index f6eff62a9d..3f10a5ced6 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcOpenGaussIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcOpenGaussIT.java
@@ -23,6 +23,7 @@ import 
org.apache.seatunnel.shade.org.apache.commons.lang3.tuple.Pair;
 
 import org.apache.seatunnel.api.table.catalog.Catalog;
 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.PrimaryKey;
 import org.apache.seatunnel.api.table.catalog.TablePath;
@@ -49,10 +50,12 @@ import java.time.Duration;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.UUID;
+import java.util.stream.Collectors;
 
 @Slf4j
 public class JdbcOpenGaussIT extends AbstractJdbcIT {
@@ -69,6 +72,13 @@ public class JdbcOpenGaussIT extends AbstractJdbcIT {
     private static final String SOURCE_TABLE = "gs_e2e_source_table";
     private static final String SINK_TABLE = "gs_e2e_sink_table";
     private static final String CATALOG_TABLE = "e2e_table_catalog";
+
+    /**
+     * OpenGauss table used to verify that catalog discovery never exposes 
dropped-column metadata
+     * placeholders as physical fields.
+     */
+    private static final String DROPPED_COLUMN_TABLE = 
"gs_catalog_dropped_column_test";
+
     private static final Integer GEN_ROWS = 100;
     private static final List<String> CONFIG_FILE =
             Lists.newArrayList("/jdbc_opengauss_source_and_sink.conf");
@@ -154,6 +164,39 @@ public class JdbcOpenGaussIT extends AbstractJdbcIT {
         Assertions.assertFalse(catalog.tableExists(targetTablePath));
     }
 
+    /**
+     * Verifies catalog discovery excludes OpenGauss metadata placeholders 
while preserving the
+     * physical order of all remaining columns.
+     */
+    @Test
+    public void testCatalogExcludesDroppedColumns() {
+        OpenGaussCatalog openGaussCatalog = (OpenGaussCatalog) catalog;
+        TablePath tablePath = TablePath.of(DATABASE, SCHEMA, 
DROPPED_COLUMN_TABLE);
+        openGaussCatalog.dropTable(tablePath, true);
+        try {
+            openGaussCatalog.executeSql(
+                    tablePath,
+                    "CREATE TABLE "
+                            + SCHEMA
+                            + "."
+                            + DROPPED_COLUMN_TABLE
+                            + " (c1 INT, c2 VARCHAR(50), c3 INT, c4 TEXT)");
+            openGaussCatalog.executeSql(
+                    tablePath,
+                    "ALTER TABLE " + SCHEMA + "." + DROPPED_COLUMN_TABLE + " 
DROP COLUMN c2");
+
+            CatalogTable table = openGaussCatalog.getTable(tablePath);
+            List<String> actualColumns =
+                    table.getTableSchema().getColumns().stream()
+                            .map(Column::getName)
+                            .collect(Collectors.toList());
+
+            Assertions.assertEquals(Arrays.asList("c1", "c3", "c4"), 
actualColumns);
+        } finally {
+            openGaussCatalog.dropTable(tablePath, true);
+        }
+    }
+
     @Test
     public void testCreateIndex() {
         String schema = "public";

Reply via email to