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;
+        }
+    }
 }

Reply via email to