This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11399-4abe9e71e14267646600dddbab860e158c4a7900 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 446ff6266886fe1406b2abc2a3553b2ea070a21a Author: Jast <[email protected]> AuthorDate: Mon Sep 28 12:07:52 2026 +0000 [Fix][Connector-V2] Filter Postgres CDC database discovery (#11399) Co-authored-by: zhangshenghang <[email protected]> Co-authored-by: davidzollo <[email protected]> Co-authored-by: jast <[email protected]> --- .../config/PostgresSourceConfigFactory.java | 5 ++ .../cdc/postgres/source/PostgresDialect.java | 8 +- .../cdc/postgres/utils/TableDiscoveryUtils.java | 30 ++++++- .../config/PostgresSourceConfigFactoryTest.java | 4 + .../postgres/utils/TableDiscoveryUtilsTest.java | 96 ++++++++++++++++++++-- 5 files changed, 136 insertions(+), 7 deletions(-) diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactory.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactory.java index 1aec7df2d5..0a03e247e1 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactory.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactory.java @@ -70,6 +70,11 @@ public class PostgresSourceConfigFactory extends JdbcSourceConfigFactory { props.setProperty("database.password", checkNotNull(password)); props.setProperty("database.port", String.valueOf(port)); props.setProperty("database.dbname", checkNotNull(databaseList.get(0))); + // Deliberately do NOT set "database.include.list": Debezium folds it into + // dataCollectionFilter() as a predicate on TableId#catalog, but PostgreSQL table ids are + // catalog-less (see PostgresSchema#readTableSchema and the event dispatchers), so every + // snapshot schema read and streaming event would be filtered out. Database scoping for + // discovery is applied in TableDiscoveryUtils#listTables via an explicit predicate. props.setProperty("plugin.name", decodingPluginName); props.setProperty("slot.name", slotName); diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresDialect.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresDialect.java index 650b6068d1..09bf8e6b60 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresDialect.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresDialect.java @@ -51,6 +51,7 @@ import io.debezium.relational.history.TableChanges; import java.sql.SQLException; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; @@ -113,9 +114,14 @@ public class PostgresDialect implements JdbcDataSourceDialect { public List<TableId> discoverDataCollections(JdbcSourceConfig sourceConfig) { PostgresSourceConfig postgresSourceConfig = (PostgresSourceConfig) sourceConfig; try (JdbcConnection jdbcConnection = openJdbcConnection(sourceConfig)) { + // Scope discovery to the configured databases via an explicit predicate instead of + // Debezium's "database.include.list", which would filter out the catalog-less + // TableIds the PostgreSQL connector uses outside of discovery. List<TableId> tables = TableDiscoveryUtils.listTables( - jdbcConnection, postgresSourceConfig.getTableFilters()); + jdbcConnection, + postgresSourceConfig.getTableFilters(), + new HashSet<>(postgresSourceConfig.getDatabaseList())::contains); this.checkAllTablesEnabledCapture(jdbcConnection, tables); return tables; } catch (SQLException e) { diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtils.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtils.java index 2b34c90bb0..e47f7eaff7 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtils.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtils.java @@ -29,12 +29,36 @@ import java.util.ArrayList; import java.util.LinkedHashSet; import java.util.List; import java.util.Set; +import java.util.function.Predicate; public class TableDiscoveryUtils { private static final Logger LOG = LoggerFactory.getLogger(TableDiscoveryUtils.class); + /** + * Reads the captured table ids from every database hosted by the connected PostgreSQL instance. + * + * <p>Filtering happens in two stages. First, {@code databaseFilter} is consulted per database + * before any metadata query is issued: PostgreSQL cannot read another database's {@code + * INFORMATION_SCHEMA} over the discovery connection, so probing foreign databases only produces + * a warning per database (see <a + * href="https://github.com/apache/seatunnel/issues/8184">#8184</a>). Second, tables read from + * an allowed database are passed through {@code tableFilters.dataCollectionFilter()}, which + * decides capture at table level. + * + * <p>The database predicate is deliberately kept outside the Debezium configuration: folding it + * into {@code database.include.list} would make {@code dataCollectionFilter()} reject the + * catalog-less {@link TableId}s used throughout the PostgreSQL connector. + * + * @param jdbc open connection to the database to discover + * @param tableFilters table-level capture filter built from the connector config + * @param databaseFilter predicate deciding which databases may be probed for tables + * @return the deduplicated table ids eligible for capture, in discovery order + */ @SuppressWarnings("MagicNumber") - public static List<TableId> listTables(JdbcConnection jdbc, RelationalTableFilters tableFilters) + public static List<TableId> listTables( + JdbcConnection jdbc, + RelationalTableFilters tableFilters, + Predicate<String> databaseFilter) throws SQLException { // Use a LinkedHashSet to deduplicate table ids. Some PostgreSQL-compatible databases // (e.g. HighGo) return the same physical table several times from @@ -67,6 +91,10 @@ public class TableDiscoveryUtils { // database and not taking them from the user ... LOG.info("Read list of available tables in each database"); for (String dbName : databaseNames) { + if (!databaseFilter.test(dbName)) { + LOG.debug("\t database '{}' is filtered out of capturing", dbName); + continue; + } try { jdbc.query( "SELECT * FROM \"" diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactoryTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactoryTest.java index 67324c0a2e..b531c60f12 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactoryTest.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/config/PostgresSourceConfigFactoryTest.java @@ -40,5 +40,9 @@ public class PostgresSourceConfigFactoryTest { Assertions.assertEquals( "never", configFactory.create(0).getDbzConfiguration().getString("snapshot.mode")); + // "database.include.list" must stay unset: Debezium turns it into a catalog predicate on + // dataCollectionFilter(), which rejects the catalog-less TableIds used by PostgreSQL. + Assertions.assertNull( + configFactory.create(0).getDbzConfiguration().getString("database.include.list")); } } diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtilsTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtilsTest.java index f68724951d..b609c3f314 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtilsTest.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/utils/TableDiscoveryUtilsTest.java @@ -22,6 +22,7 @@ import org.apache.seatunnel.connectors.seatunnel.cdc.postgres.config.PostgresSou import org.junit.jupiter.api.Test; +import io.debezium.config.Configuration; import io.debezium.connector.postgresql.connection.PostgresConnection; import io.debezium.jdbc.JdbcConfiguration; import io.debezium.jdbc.JdbcConnection; @@ -33,10 +34,14 @@ import java.sql.SQLException; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Predicate; +import java.util.stream.Collectors; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -78,14 +83,18 @@ public class TableDiscoveryUtilsTest { return resultSet; } - private static RelationalTableFilters tableFilters() { + private static Predicate<String> testdbFilter() { + return new HashSet<>(Collections.singletonList("testdb"))::contains; + } + + private static RelationalTableFilters tableFilters(String database) { PostgresSourceConfig config = (PostgresSourceConfig) new PostgresSourceConfigFactory() .hostname("localhost") .username("user") .password("password") - .databaseList("testdb") + .databaseList(database) .create(0); return config.getTableFilters(); } @@ -103,7 +112,8 @@ public class TableDiscoveryUtilsTest { } List<TableId> tableIds = - TableDiscoveryUtils.listTables(new FakePostgresConnection(rows), tableFilters()); + TableDiscoveryUtils.listTables( + new FakePostgresConnection(rows), tableFilters("testdb"), testdbFilter()); assertEquals(1, tableIds.size()); assertEquals(new TableId("highgo", "testdb", "test_a"), tableIds.get(0)); @@ -119,7 +129,8 @@ public class TableDiscoveryUtilsTest { new String[] {"testdb", "inventory", "shipments"}); List<TableId> tableIds = - TableDiscoveryUtils.listTables(new FakePostgresConnection(rows), tableFilters()); + TableDiscoveryUtils.listTables( + new FakePostgresConnection(rows), tableFilters("testdb"), testdbFilter()); assertEquals( Arrays.asList( @@ -141,7 +152,8 @@ public class TableDiscoveryUtilsTest { new String[] {"testdb", "public", "orders"}); List<TableId> tableIds = - TableDiscoveryUtils.listTables(new FakePostgresConnection(rows), tableFilters()); + TableDiscoveryUtils.listTables( + new FakePostgresConnection(rows), tableFilters("testdb"), testdbFilter()); assertEquals( Arrays.asList( @@ -150,4 +162,78 @@ public class TableDiscoveryUtilsTest { new TableId("testdb", "public", "users")), tableIds); } + + /** The real Debezium filter must keep accepting the catalog-less TableIds used at runtime. */ + @Test + public void shouldKeepAcceptingCatalogLessTableIds() { + assertTrue( + tableFilters("testdb") + .dataCollectionFilter() + .isIncluded(new TableId(null, "public", "orders"))); + assertTrue( + tableFilters("testdb") + .dataCollectionFilter() + .isIncluded(new TableId("highgo", "testdb", "test_a"))); + } + + /** Only databases accepted by the explicit database predicate are probed for tables. */ + @Test + public void shouldOnlyQueryDatabasesAllowedByConfiguredFilter() throws SQLException { + RelationalTableFilters tableFilters = tableFilters("selected"); + Predicate<String> databaseFilter = + new HashSet<>(Collections.singletonList("selected"))::contains; + + try (MockJdbcConnection jdbc = new MockJdbcConnection()) { + List<TableId> tableIds = + TableDiscoveryUtils.listTables(jdbc, tableFilters, databaseFilter); + + assertEquals( + Collections.singletonList(new TableId("selected", "public", "orders")), + tableIds); + List<String> metadataQueries = + jdbc.getQueries().stream() + .filter(query -> query.contains("INFORMATION_SCHEMA.TABLES")) + .collect(Collectors.toList()); + assertEquals(1, metadataQueries.size()); + assertTrue(metadataQueries.get(0).contains("\"selected\"")); + assertTrue( + jdbc.getQueries().stream().noneMatch(query -> query.contains("\"unwanted\""))); + } + } + + private static class MockJdbcConnection extends JdbcConnection { + private final List<String> queries = new ArrayList<>(); + + MockJdbcConnection() { + super( + JdbcConfiguration.adapt(Configuration.from(Collections.emptyMap())), + config -> null, + "\"", + "\""); + } + + @Override + public JdbcConnection query(String query, ResultSetConsumer resultConsumer) + throws SQLException { + queries.add(query); + ResultSet resultSet = mock(ResultSet.class); + if (query.equals("select datname from pg_database")) { + when(resultSet.next()).thenReturn(true, true, false); + when(resultSet.getString(1)).thenReturn("selected", "unwanted"); + } else if (query.contains("\"selected\"")) { + when(resultSet.next()).thenReturn(true, false); + when(resultSet.getString(1)).thenReturn("selected"); + when(resultSet.getString(2)).thenReturn("public"); + when(resultSet.getString(3)).thenReturn("orders"); + } else { + throw new AssertionError("Unexpected database query: " + query); + } + resultConsumer.accept(resultSet); + return this; + } + + List<String> getQueries() { + return queries; + } + } }
