eskabetxe commented on PR #209: URL: https://github.com/apache/flink-connector-jdbc/pull/209#issuecomment-5886517953
> Two implementation notes from building this locally. The design points are on the [DISCUSS] thread. > > **A race in `DatabaseSplitterEnumeratorTest.testStartExportsGlobalSnapshotSynchronouslyBeforeAnySplit`.** The test checks `createGlobalSnapshotCount == 1` right after `start()`. The background worker starts one table splitter per table, and `TableSplitterEnumerator.start()` calls `createGlobalSnapshot()` again. On the real provider that is a no-op, since it returns early when `snapshotId != null`. But `FakeConnectionProvider` only increments a counter, and its `newInstance()` returns `this`, so the counter keeps going up. > > I drove the real `DatabaseSplitterEnumerator` with that same fake, 20 times: the counter was 1 immediately after `start()` in all 20 runs, and reached 2 about 20ms later in every run. So the assertion holds only while the worker has not started the table splitter yet. On a loaded CI machine that window can close first. If the fake did nothing when it already holds a snapshot, like the real provider does, the race would be gone and the fake would behave closer to production. > > **The pool size and splitter concurrency are only equal by chance.** `MAX_CONCURRENT_TABLE_SPLITTERS = 4` in `DatabaseSplitterEnumerator` happens to match `DEFAULT_POOL_SIZE = 4` in `AbstractConnectionProvider`, and the comment about never exceeding pool capacity relies on that. It is true today, and finished splitters do get closed and return their connection before `fillActiveTableSplitters()` runs again. But `maxPoolSize()` can be overridden by a dialect, and the two constants live in different classes with nothing linking them. If a dialect used a smaller pool, the fill loop would block for `connectionCheckTimeout` (60s by default) and then fail. Deriving one from the other, or making the pool size a `ConnectionOptions` setting, would make it explicit. > > Minor: in `JdbcSourceSplitSerializer.deserializeJdbcSourceSplit` the return value of `in.read(parametersBytes)` is ignored. It is pre-existing, and the same in the v0 path, but `in.readFully(...)` would be safer. Thanks for the detailed review @SEPURI-SAI-KRISHNA, - FakeConnectionProvider counter — fixed: createGlobalSnapshot() is now idempotent (compareAndSet(0, 1)), mirroring the real provider's "an existing exported snapshot is kept" behavior, so derived table splitters re-calling it after newInstance() no longer inflate the count. - Fan-out vs pool size — agreed, the two independent 4s were a bug waiting to happen. Reworked as follows: ConnectionProvider gains an abstract int maxPoolSize(), which AbstractConnectionProvider implements as jdbcOptions.getConnectionPoolSize(); ConnectionOptionsBuilder gains withConnectionPoolSize(int) (default 4). DatabaseSplitterEnumerator derives its fan-out solely from maxPoolSize() (MAX_CONCURRENT_TABLE_SPLITTERS is removed), and the Postgres dialect's private POOL_SIZE override is gone so the option actually controls the pool. A Math.max(1, ...) clamp prevents a custom provider reporting a non-positive capacity from deadlocking enumeration. The interface default deliberately isn't a constant — a default couldn't see the configured options and could silently disagree with the real pool, so implementors must report the value they honor. New tests assert fan-out for pool sizes 1/4/16/0. - in.read return value — fixed: both deserialization paths in JdbcSourceSplitSerializer now use readFully(...). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
