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]