nzw921rx commented on PR #11512:
URL: https://github.com/apache/seatunnel/pull/11512#issuecomment-5068763565

   > Thanks @nzw921rx I considered this again.
   > 
   > The CDC scope is already explicit in CdcEnumeratorProgressProvider and 
CdcEnumeratorProgressReport. Provider here represents a local non-blocking 
capability that returns the latest connector progress snapshot. I kept 
Observability out of the type name because these types describe connector-owned 
progress facts rather than the engine observability service itself.
   > 
   > I have kept the current names for now, but can revise them if you think a 
different name would better match the project convention. Let me know
   
   
   Thank you for your reply. My idea was generated by AI and I would like to 
discuss it with you. It's an honor
   
   ## 1. The reporting mechanism is duplicated across nearly every layer
   
   The current implementation introduces two parallel reporting paths:
   
   * `CdcReaderProgressProvider` / `CdcEnumeratorProgressProvider`
   * `CdcReaderProgressReport` / `CdcEnumeratorProgressReport`
   * reader/enumerator task methods for collection, source vertex lookup, and 
sequence generation
   * `CdcReaderProgressEnvelope` / `CdcEnumeratorProgressEnvelope`
   * separate reader/enumerator lists in `ReportCdcProgressOperation`
   * separate serialization branches
   * separate maps, update methods, and ordering methods in `CdcProgressService`
   
   Reader and enumerator reports should remain different because they own 
different facts. However, provider discovery, engine identity, 
execution-attempt ordering, transport, serialization, and latest-report storage 
are common reporting infrastructure and should not be duplicated.
   
   A shared connector-facing abstraction could look like:
   
   ```java
   public interface CdcProgressReport extends Serializable {}
   
   public interface CdcProgressProvider<R extends CdcProgressReport> {
       R getCdcProgress();
   }
   ```
   
   The engine could then use one tagged envelope and one ordering/storage path 
while preserving reader- and enumerator-specific report payloads.
   
   Without a common top-level protocol, every additional report owner or 
progress dimension is likely to introduce another provider, envelope, 
serializer branch, task collection branch, and state map.
   
   ## 2. Reader and enumerator lifecycle ownership is conflated
   
   `CdcProgressLifecycle` is shared by both reader and enumerator reports, but 
`SnapshotSplitAssigner` reports `CATCH_UP` once split assignment is complete:
   
   ```java
   assignerCompleted
           ? CdcProgressLifecycle.CATCH_UP
           : CdcProgressLifecycle.SNAPSHOT
   ```
   
   `CATCH_UP` describes reader consumption relative to the snapshot high 
watermark. The enumerator does not own or directly observe that state, and 
assignment completion does not prove that readers are still catching up. 
Readers may already be incremental, or different readers may temporarily be in 
different states.
   
   The enumerator should report assignment-specific state, such as whether 
snapshot discovery or assignment is still running. The source-level CDC 
lifecycle should be derived by the engine from reader reports without hiding 
mixed reader states.
   
   This is also important for the STIP-30 principle that each component reports 
only facts it owns.
   
   ## 3. Enumerator collection does not follow the agreed STIP-30 path
   
   STIP-30 specifies two distinct collection paths:
   
   * reader reports are periodically sampled on workers and transported to the 
coordinator;
   * enumerator progress is collected locally on the coordinator and is not 
sampled through the worker task loop.
   
   The current implementation scans both `SourceSeaTunnelTask` and 
`SourceSplitEnumeratorTask` from `TaskExecutionService.collectCdcProgress()` 
and places both report types into `ReportCdcProgressOperation`.
   
   Even if the enumerator task currently runs on the coordinator, this still 
routes coordinator-owned state through the generic task sampler and operation 
transport instead of collecting it directly through the coordinator-side 
progress service. It introduces unnecessary transport and diverges from the 
agreed design.
   
   ## 4. Engine transport constraints leak into the connector-facing API
   
   All public progress DTOs implement `Serializable`, and `CdcProgressValue` 
requires:
   
   ```java
   T extends Serializable
   ```
   
   This makes an engine transport requirement part of the connector-facing SPI. 
The engine already has explicit serialization code, so the public report model 
should not be shaped around the current Hazelcast transport implementation 
unless that restriction is intentionally part of the connector contract.
   
   ## Suggested direction
   
   I suggest keeping the reader and enumerator payload models separate, while 
consolidating the common mechanics:
   
   1. Introduce one common `CdcProgressProvider<R>` and report marker.
   2. Introduce one engine-internal report source abstraction for task 
collection.
   3. Use one envelope containing owner type, task identity, execution attempt, 
sequence, observation time, and report payload.
   4. Use one latest-report ordering and storage implementation.
   5. Keep enumerator assignment state separate from reader CDC lifecycle.
   6. Collect enumerator state directly on the coordinator as specified by 
STIP-30.
   
   Because this is a foundational connector-facing and engine-level contract, 
it will become significantly harder to consolidate after checkpoint, restore, 
aggregation, REST, and additional CDC connectors depend on the current parallel 
hierarchy. I think this should be addressed before merge.
   


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