goutamadwant opened a new pull request, #11512:
URL: https://github.com/apache/seatunnel/pull/11512

   ### Purpose of this pull request
   
   Implement the first experimental runtime reporting slice of STIP-30 
(#11364), following the lifecycle discussion in #11395. This is an 
engine-internal foundation, not a REST/UI feature or a change to checkpoint 
semantics.
   
   - Keep reader lifecycle and consumed-position facts separate from enumerator 
discovery/assignment facts.
   - Use immutable connector-owned reports, a tagged engine envelope, an 
explicit wire codec with collection-count limits and coordinator-local 
latest-only storage. These limits do not bound individual strings or total 
frame bytes.
   - Opt sources in through `SupportCdcProgress`; filter non-CDC sources before 
coordinator slot lookups.
   - Capture detached reader coordinates only after successful emission. 
Configured/restored starts are `BEST_EFFORT`, not proof of consumption. A 
mutable connector offset, failed deserialization or failed collection cannot 
advance a published position.
   - Publish enumerator counts and bounded details as one completed immutable 
snapshot without making polling wait for chunk generation or source I/O. 
Include snapshot-only provider wiring.
   - Isolate failing providers, preserve subsequent monitoring ticks and allow 
at most one outstanding enumerator request per worker.
   - Close pipeline observation scopes on terminal cleanup, including savepoint 
completion. Fence old coordinator initialization/callbacks by generation and 
roll back failed initialization using JobMaster ownership. Delayed metrics 
cleanup no longer performs an unowned CDC removal that could erase a 
replacement scope.
   - Preserve the previous reader constructor overload and existing engine 
identity serialization. No dependency, connector registration or 
checkpoint-format change is introduced.
   
   ### Runtime reporting flow
   
   Reader: successful record emission -> detached immutable tracker state -> 
worker sampling -> batched operation -> active coordinator latest-only store.
   
   Enumerator: completed assigner mutation -> immutable assignment snapshot -> 
coordinator identifies opted-in source tasks in running plans and assigned 
workers -> asynchronous worker collection -> generation-checked publication.
   
   The coordinator resolves current placement; this is not a separate 
persistent enumerator registry. Remote placement does not transfer polling 
ownership to the worker. One failed provider or worker must not prevent other 
completed reports or later ticks.
   
   Per-owner attempt and sequence ordering is relative to the report already 
stored. It is **not** validation against an authoritative expected deployment: 
an empty reopened scope can receive an older reader observation before its 
first newer report. Generation fencing protects captured coordinator work, not 
every delayed network packet. Reader observations are not an atomic distributed 
snapshot.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes: experimental provider/report interfaces and internal runtime 
collection, with EN/ZH developer documentation. Existing job configuration and 
output are unchanged. Reports are not exposed through REST, CLI or metrics in 
this slice.
   
   Completed-checkpoint and restored-position facts remain `UNSUPPORTED` until 
their actual callbacks are connected. Split assignment does not prove restore 
origin. Native positions must contain coordinates only, never credentials or 
connection URLs. Each fact carries its own accuracy; an unchanged position does 
not prove lag or a stalled source.
   
   The immutable reader publication adds real per-emission allocation. A 
retained synthetic single-JVM measurement on JDK 11 observed approximately 0.8 
microseconds and 1,136 bytes per publication for eight coordinates; 64 
coordinates cost approximately 5.3 microseconds and 7,008 bytes. These include 
conversion/publication and clock work, are not controlled JMH or 
production-throughput results, and exclude real token construction and 
downstream work. No zero-overhead or production-performance acceptance claim is 
made. The previous lower-allocation implementation was incorrect because it 
retained a mutable offset.
   
   ### How was this patch tested?
   
   Latest follow-up verification:
   
   - 409 API tests passed on JDK 11.
   - 122 CDC base tests passed on JDK 11 and JDK 8 runtime, including opt-in 
allocation measurements, lazy-chunking concurrency, snapshot-only wiring and 
immutable reader publication.
   - Four MongoDB component tests passed on each runtime using the actual 
emitter and `ChangeStreamOffset`, covering deserializer failure, collector 
failure and polling during a latch-held emission. Before the fix, all three new 
emitter regressions failed; the tracker regression independently reproduced 
three failures.
   - JDK 8 executed Java-8-targeted classes compiled with JDK 11; this 
follow-up does not claim a separate clean JDK 8 compilation. The MongoDB 
lifecycle test resolved the current base module classes, not a stale installed 
artifact.
   - Four Markdown checks passed for the revised EN/ZH contract.
   - The final JDK 11 focused engine selection reported 112 cases: 110 passed, 
two existing skips and zero failures/errors, including the reproduced 
stale-generation restore and delayed-cleanup regressions.
   - The final JDK 8 runtime engine selection passed all 32 cases, covering 
coordinator lifecycle, delayed cleanup, progress storage, selected 
JobMaster/savepoint behavior and serialization.
   - The full JDK 11 engine-server unit-module run reported 594 cases: 588 
passed, five skipped and one error in 
`SavePointTest.testSavePointButJobGoingToFail`. The observed job stayed RUNNING 
instead of reaching FAILED within the existing timeout. The isolated method 
then passed on both the patched tree and a clean checkout of the pre-fix 
commit. The initiating full-suite cause remains unproven; this is not a green 
full-module result or proof that the failure is unrelated.
   
   The regressions cover pre-initialization null handling, snapshot 
consistency, provider exception isolation, non-CDC filtering, 
outstanding-request control, malformed/oversized codec input, 
terminal/savepoint cleanup, failed initialization, owner replacement and 
coordinator generation changes. Existing serialization registration and 
checkpoint state remain unchanged.
   
   Commands (select the JDK through JAVA_HOME; module prerequisites must be 
available):
   
   ```shell
   ./mvnw -B -T 1 -pl seatunnel-api test
   ./mvnw -B -T 1 \
     -pl seatunnel-connectors-v2/connector-cdc/connector-cdc-base \
     -Dcdc.progress.benchmark=true verify
   ./mvnw -B -T 1 \
     -pl 
seatunnel-connectors-v2/connector-cdc/connector-cdc-base,seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb
 \
     -Dtest=MongoDBRecordEmitterProgressTest,ChangeStreamOffsetTest \
     -Dsurefire.failIfNoSpecifiedTests=false test
   ./mvnw -B -T 1 -pl seatunnel-engine/seatunnel-engine-server test
   ```
   
   These are focused/module checks, not a claim of full-reactor, live-database 
recovery or aggregate CI success. The dependent real MySQL restore/progress 
slice remains draft #12137: after this foundation lands, it must be rebased to 
a test-only diff and executed against the final contract. Synthetic allocation 
data does not close end-to-end performance acceptance.
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
   * [x] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR.
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     2. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
     3. Add ci label in 
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
     4. Add e2e testcase in 
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
     5. Update connector 
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
   
   No new connector, binary dependency, plugin mapping or incompatible 
configuration change is introduced. Real database E2E remains the separate 
dependent PR, not completed by this slice's component tests.
   


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