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";