Hi Hongshun, Thanks for picking this up! wglbd12 is actually my JIRA account, so the ticket is already with the right person.
The fix is up for review: https://github.com/apache/flink-cdc/pull/4534 It restores the per-table lifetime of the splitter connection (as it was before the async splitting change): the connection is acquired lazily and released once a table has been split, and close() prevents it from being re-acquired afterwards. Unit tests cover the evenly and unevenly sized paths. I also verified it on a Flink 1.20 session cluster with 30 postgres-cdc jobs against one database and the default pool size: all jobs start, and JobManager connections drop from 31 to 1. The CI workflows are waiting for approval since this is my first contribution to the repo. Could you approve them and take a look when you have time? Best, Kyrie Hongshun Wang <[email protected]> 于2026年9月17日周四 10:56写道: > > Hi Kyrie > wglbd12 is willing to fix it[1] ,and I assign this ticket to him. Maybe you > can talk with him > > Best, > Hongshun > > [1] https://issues.apache.org/jira/browse/FLINK-40623 > > > On 2026/09/10 02:47:47 Kyrie wrote: > > Hi Flink CDC devs, > > > > I'd like to report a bug in the CDC incremental source framework > > (flink-cdc-base) that I hit while upgrading our PostgreSQL CDC jobs > > from 3.2.1 to 3.6.0. I have requested an ASF JIRA account but it is > > not approved yet, so I'm posting the full analysis here first. If a > > committer prefers to open the ticket on my behalf, please feel free > > to; otherwise I will file it as soon as the account is approved > > (component: Flink CDC). > > > > == Summary == > > > > Since FLINK-34688 (asynchronous chunk splitting, PR #3510, included > > from cdc-3.3.0 on), JdbcSourceChunkSplitter opens its JDBC connection > > in open() and only closes it in close(). SnapshotSplitAssigner calls > > chunkSplitter.close() only from its own close(), i.e. when the job > > terminates, not when splitting is finished. The connection comes from > > the HikariCP pool in JdbcConnectionPools, which is a JVM-wide > > singleton shared by every job on the same JobManager that targets the > > same host/port/user/database. The pool size is connection.pool.size > > (default 20) and the borrow timeout is connect.timeout (default 30 s). > > > > Net effect: every running incremental-snapshot job pins one pooled > > connection on the JobManager for its whole lifetime. Once 20 such jobs > > are running against the same database, the 21st job fails in > > createEnumerator -> discoverDataCollections() with a Hikari timeout > > and loops in RESTARTING forever. Nothing in the message points at > > connection.pool.size. > > > > == Environment == > > > > * Flink 1.20.0 standalone session cluster (1 JM + 1 TM), JDK 17 > > * flink-sql-connector-postgres-cdc 3.6.0-1.20 (Debezium 1.9.8.Final) > > * PostgreSQL 15 (cloud RDS), 30 independent Flink SQL jobs, one > > postgres-cdc table each, all on the same host/port/user/database > > * scan.incremental.snapshot.enabled = true, scan.startup.mode = > > initial, connection.pool.size not set (default 20) > > * 3.2.1 with the same DDL runs all 30 jobs fine; 3.2.1 still used > > try-with-resources per generateSplits(). > > > > == Symptom == > > > > Jobs 1..20 start normally. From job 21 on, every new job fails during > > enumerator creation: > > > > Caused by: org.apache.flink.util.FlinkRuntimeException: Failed to > > discover captured tables for enumerator > > Caused by: org.apache.flink.util.FlinkRuntimeException: > > java.sql.SQLTransientConnectionException: connection-pool-<host>:5432 > > - Connection is not available, request timed out after 30000ms. > > at > > org.apache.flink.cdc.connectors.base.relational.connection.JdbcConnectionFactory.connect(...) > > at > > org.apache.flink.cdc.connectors.postgres.source.PostgresDialect.openJdbcConnection(...) > > at > > org.apache.flink.cdc.connectors.postgres.source.PostgresDialect.discoverDataCollections(...) > > at > > org.apache.flink.cdc.connectors.postgres.source.PostgresSourceBuilder$PostgresIncrementalSource.createEnumerator(...) > > > > * Submitting the failing job alone (after cancelling the other failing > > ones) does not help. > > * Restarting the TaskManager does not help (the pool that is exhausted > > lives in the JobManager). > > * Restarting the JobManager "fixes" it until 20 jobs are running again. > > * pg_stat_activity shows exactly 20 idle backends from the JobManager, > > one per running job, created at each job's start and never used again. > > Their last statement is the chunk-split query: > > SELECT MAX("id") FROM (SELECT "id" FROM "<db>"."public"."<table>" > > WHERE "id" >= $1 ORDER BY "id" ASC LIMIT 8096) AS T > > > > == Code references (release-3.6) == > > > > * flink-cdc-base .../source/assigner/splitter/JdbcSourceChunkSplitter.java > > open(): this.jdbcConnection = dialect.openJdbcConnection(sourceConfig); > > close(): if (jdbcConnection != null) { jdbcConnection.close(); } > > release-3.2 had: try (JdbcConnection jdbc = > > dialect.openJdbcConnection(sourceConfig)) { ... } inside > > generateSplits(). > > * flink-cdc-base .../source/assigner/SnapshotSplitAssigner.java > > open() -> chunkSplitter.open(); ... startAsynchronouslySplit(); > > close() -> chunkSplitter.close(); // the only place it is closed > > onFinishedSplits() flips assignerStatus to finished but does not > > close the splitter. > > * flink-cdc-base .../relational/connection/JdbcConnectionFactory.java > > > > JdbcConnectionPools.getInstance(factory).getOrCreateConnectionPool(poolId, > > sourceConfig).getConnection() > > JdbcConnectionPools is a singleton; the pool is keyed per > > host/port/user/database, i.e. shared across jobs in the JVM. > > * flink-cdc-base .../relational/connection/JdbcConnectionPoolFactory.java > > config.setMaximumPoolSize(sourceConfig.getConnectionPoolSize()); > > // default 20 > > > > config.setConnectionTimeout(sourceConfig.getConnectTimeout().toMillis()); > > // default 30 s > > > > == Expected behaviour / proposal == > > > > The chunk splitter should give its connection back once splitting for > > all remaining tables is done (e.g. close it at the end of > > splitChunksForRemainingTables(), or go back to a try-with-resources > > per generateSplits() call, re-opening lazily if newly added tables > > arrive). Then the number of concurrently running jobs is no longer > > bounded by connection.pool.size on the JobManager. At minimum the > > limitation should be documented and the timeout should be reported > > with a hint about connection.pool.size. > > > > == Workaround == > > > > Set 'connection.pool.size' in the table DDL to a value larger than the > > number of jobs sharing the JobManager (we use 64 for 30 jobs) and > > restart the JobManager, because the pool size is fixed when the pool > > is first created. > > > > == How to reproduce == > > > > 1. Session cluster, postgres-cdc 3.6.0, incremental snapshot enabled, > > connection.pool.size at default. > > 2. Submit 21 independent SQL jobs, each reading one table from the > > same PostgreSQL database. > > 3. Job 21 fails in createEnumerator with the exception above; > > pg_stat_activity shows 20 idle JobManager backends whose last query is > > the chunk-split SELECT MAX(...) LIMIT 8096. > > > > I'm happy to work on a fix once the JIRA exists and is assigned. > > > > Thanks. > >
