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]

Reply via email to