This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] 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 446ff62668 [Fix][Connector-V2] Filter Postgres CDC database discovery
(#11399)
446ff62668 is described below
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;
+ }
+ }
}