DanielLeens commented on PR #11932:
URL: https://github.com/apache/seatunnel/pull/11932#issuecomment-5476988683

   Thanks for reviewing, @JeremyXin. I re-checked the current head (`4414b55b`, 
unchanged since 2026-08-22 — no new commits since my last comment) against your 
two specific asks plus your overall read, since the head hasn't moved and the 
outstanding items from mine and @SEZ9's earlier discussion are still all 
present in the code.
   
   **On your two points:**
   
   1. **Impact scope of `Collector.restoreSchema`**: confirmed via the 
instantiation site — `IncrementalSourceReader` is only ever constructed in 
`IncrementalSource.java` (single call site), and that base class is extended by 
`Db2IncrementalSource`, `OracleIncrementalSource`, `MySqlIncrementalSource`, 
`PostgresIncrementalSource`, and `SqlServerIncrementalSource`. So this is the 
shared reader for all five Debezium-based CDC connectors, not a MySQL-only path 
— good for coverage, but it also means the still-open gaps below (Issues 1 and 
2) affect all five, not just MySQL.
   
   2. **`restoredCheckpointTables` read+clear atomicity**: I traced the actual 
call paths rather than just checking for a `synchronized`/lock keyword, and the 
"single thread" assumption doesn't hold. The write side — `initializedState()` 
setting `restoredCheckpointTables` (`IncrementalSourceReader.java:239`) — is 
reached via `addSplits()` ← `SourceFlowLifeCycle.receivedSplits()` ← 
`SourceSeaTunnelTask.receivedSourceSplit()` ← 
`AssignSplitOperation.runInternal()`, and `AssignSplitOperation extends 
TracingOperation implements IdentifiedDataSerializable` — a Hazelcast 
`Operation`, executed on a Hazelcast operation-execution thread. The read side 
— `restoreCollectorSchema()` (`:126-133`) — runs from `pollNext()` ← 
`SourceFlowLifeCycle.collect()`, which is the task's own execution-loop thread. 
These are two different threads. `restoredCheckpointTables` being `volatile` 
gives visibility but not atomicity of the "read, then null out" sequence, so 
there's a real (if narrow)
  window: if a newly-assigned split's `initializedState()` sets 
`restoredCheckpointTables = Y` while `pollNext()` is mid-way through reading 
the previous value and about to null it, `Y` can get silently wiped before 
`output.restoreSchema(Y)` ever runs for it — a dropped schema restore, not just 
a theoretical concern. I'd raise this as a genuine correctness finding rather 
than something a doc comment alone resolves; it likely needs the 
assignment/consumption to be a single atomic operation (e.g. `getAndSet(null)` 
at the point of use, or moving the field through a queue), not just a note 
about single-threadedness.
   
   **On the overall read**, I want to flag a factual correction before the 
"core approach is sound" / Comment-level conclusion gets read as "close to 
mergeable": the current head hasn't fixed any of the previously confirmed 
blocking items, and one of them contradicts your summary directly. 
Specifically, at `IncrementalSourceReader.java:229-240` right now:
   ```java
   if (incrementalSplit.getCheckpointTables() != null) {
       ...
       
debeziumDeserializationSchema.restoreCheckpointProducedType(incrementalSplit.getCheckpointTables());
       // Only roll the collector back to checkpointed tables when the dialect 
can also
       // restore the deserializer's produced type from the same checkpoint 
snapshot.
       if (debeziumDeserializationSchema.getSchemaChangeResolver() != null) {
           restoredCheckpointTables = incrementalSplit.getCheckpointTables();
       }
   }
   ```
   the deserializer-side `restoreCheckpointProducedType(...)` call runs 
unconditionally, *outside* the `getSchemaChangeResolver() != null` gate that 
protects the collector-side restore — so the source and deserializer sides do 
**not** restore consistently on a resolver-less dialect; that's exactly Issue 2 
from the 2026-08-23/24 discussion, still open. Issue 1 (no `isEmpty()` guard 
alongside the `!= null` check on the same line) is also still open. And on the 
JDBC side, `JdbcSinkWriter.snapshotState()` (`:212-213`) still unconditionally 
returns `Collections.singletonList(new JdbcSinkState(null, tableSchema))` for 
every writer, not just schema-evolution-affected ones — that's Issue 5 / my 
Finding A, also still open, and it's the one that directly interacts with the 
error-sink lifecycle gate I flagged earlier.
   
   None of this is new work needed from you — just flagging that these three 
are pre-existing, already-agreed blockers (not resolved by anything in this 
thread), so my and @SEZ9's "not recommended for merge in current form" 
conclusion from 2026-08-26 still stands as of the current head. Your two new 
points (scope confirmation, and especially the concurrency finding) are 
genuinely useful additions to that list, though, and I've folded them in above. 
Agreed on the branch-conflicts point too — those need to be cleared regardless.
   
   @davidzollo — still hoping to hear back on reconciling this branch with 
#11780 before the next round of fixes; happy to do a full fresh re-review as 
soon as the open items (empty-list guard, resolver-gated deserializer restore, 
opt-in JDBC state emission, the split-assign/pollNext race above, plus Issues 
3/6/7/8 and the docs/upgrade note) land together.


-- 
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