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

   Round 9 — verification pass, no new commits since round 8. I re-fetched the 
PR: `headRefOid` is still `f243d64dadb96be96d18a1379b0f23b653bb98cb`, 
byte-identical to what round 8 reviewed two days ago. Since the source tree is 
unchanged, I did not re-derive every finding from a blank page; instead I 
independently re-verified the load-bearing claims against the real files and a 
real, fresh CI run, and I am reporting exactly what still holds and what 
changed operationally (the CI run itself was retried since round 8).
   
   **Direct answer to the open question about a "Zeta compile break":** at 
`f243d64d` the code compiles cleanly. I confirmed this by reading the full fork 
CI log for the currently-failing job rather than trusting the summary. There is 
no `error:`, no `cannot find symbol`, no `COMPILATION ERROR` anywhere in the 
log. The module `seatunnel-engine-client` reaches `mvn ... test` and fails only 
at the surefire phase, on one specific test class. Whatever compile-break 
existed at an earlier commit in this PR's history is not present at the current 
head. The CI redness at `f243d64d` is a test-fixture bug, not a build break — 
see Issue 1.
   
   # What Problem Does This PR Solve?
   
   SeaTunnel currently has no in-engine way to enrich a streaming fact source 
with a dimension table that is itself changing; today that means 
pre-materialising the dimension upstream or doing a per-row lookup in a 
transform/sink with no consistent view of the dimension across the job. This PR 
adds a first in-engine dynamic-lookup runtime ("M0"): a `dynamic_lookup { ... 
}` config block, a dedicated two-input `DynamicLookupAction` (not a Transform), 
a keyed in-memory map fed by the dimension changelog and probed by fact rows, 
and a "fact gate" that holds the Kafka fact source closed until the dimension 
side has a head start.
   
   **One-sentence summary:** a port-aware, two-input lookup-join operator added 
to the Zeta DAG, with a gated fact source and checkpointed dimension state, 
deliberately fenced to a narrow M0 feature set (`IN_MEMORY`, `ttl=NONE`, 
`APPEND_ONLY` facts, `FAIL` on PK update/schema change).
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   Critical path, verified by direct read at the current head:
   
   - 
`seatunnel-engine/seatunnel-engine-core/.../parse/MultipleTableJobConfigParser.java`
 — parses/validates `dynamic_lookup`.
   - 
`seatunnel-engine/seatunnel-engine-server/.../dag/physical/PhysicalPlanGenerator.java`
 — allocates the in-JVM queues (fact-input, dimension-input, fact-gate-command) 
and validates the pipeline shape.
   - 
`seatunnel-engine/seatunnel-engine-server/.../task/flow/DynamicLookupFlowLifeCycle.java`
 — the operator itself, single-threaded data plane.
   - 
`seatunnel-engine/seatunnel-engine-server/.../task/flow/SourceFlowLifeCycle.java`
 — gate-command draining, gate snapshot/restore.
   - 
`seatunnel-connectors-v2/connector-kafka/.../source/KafkaSourceReader.java` — 
the gate-capable Kafka reader.
   - 
`seatunnel-engine/seatunnel-engine-server/.../checkpoint/{CheckpointCoordinator,CompletedCheckpoint,CompletedCheckpointCodec,PendingCheckpoint}.java`
 — the checkpoint envelope this PR adds globally.
   
   Runtime path (traced from the real call chain, not the PR description):
   
   ```
   JOB START
     -> PhysicalPlanGenerator allocates fact-input / dimension-input / 
fact-gate-command(cap=1) queues
        these are in-JVM queues, so fact source + dimension source + lookup 
operator must share one task group
   INIT   SourceFlowLifeCycle.init() [task thread]
     -> isDynamicLookupFactGate() -> reader.prepareClosedGate(); 
KafkaSourceReader.gateOpen=false
   RESTORE  restoreState() [Hazelcast operation thread]
     -> DynamicLookupFlowLifeCycle.restoreState(): dimensionState.clear(); 
decode envelope (magic+version+len+digest)
     -> SourceFlowLifeCycle.restoreFactGateState() -> reader.restoreGateState()
     -> restoreComplete.complete(null)  (happens-before edge into open())
   OPEN   cycle.open() [task thread]
     -> if restoredFromDurableLookupState && !gateOpened -> openFactGate()
     -> Kafka: readerOpened=true; if activateOnOpen -> activateStagedSplits()
   RUNNING
     [TASK THREAD] DynamicLookupFlowLifeCycle.collect()
        drainPort(DIMENSION) then drainPort(FACT), queue.poll(10ms)/record, 
idle -> sleep(100ms)
        dimension row -> applyDimension() 
[DynamicLookupFlowLifeCycle.java:252-276]
            INSERT/UPDATE_AFTER -> dimensionState.put ; DELETE -> 
dimensionState.remove
            -> enforceDimensionStateBudget() [:276 -> :506-524], which 
re-serialises the WHOLE map every call [Issue 4]
        fact row -> projectFact()
            dimensionState.get(key); INNER + miss -> return null, row silently 
dropped, no counter/log [Issue 8]
     [TASK THREAD] SourceFlowLifeCycle.collect()
        drainGateCommands() -- OUTSIDE collector.getCheckpointLock() --
           -> reader.applyGateCommand(OPEN) -> activateStagedSplits(): gateOpen 
false->true, drains staged splits
        reader.pollNext(collector)  (lock only taken inside, per emitted record)
                       || CONCURRENT ||
     [HAZELCAST OP THREAD] CheckpointBarrierTriggerOperation -> 
SourceFlowLifeCycle.triggerBarrier()
        synchronized(collector.getCheckpointLock()) { snapshotSourceState() -> 
reader.snapshotGate() }
           snapshotGate() calls snapshotState() (acquires+releases gateLock 
once), THEN re-acquires gateLock
           separately to read gateOpen/stagedNoMoreSplits for the same 
SourceGateState [Issue 3]
           snapshotState() itself: branch chosen under gateLock, then `if 
(!gateOpen)` re-read AFTER the lock
           is released [Issue 2]
   CHECKPOINT (lookup operator) [task thread, via barrier record]
     -> alignBarrier() under checkpointStateLock (verified correct) -> 
addState(dimension envelope)
     -> arm pendingFactGateOpenCheckpointId (once)
   notifyCheckpointComplete(id) [checkpoint callback thread]
     -> if pending==id -> openFactGate() -> gateQueue.offer(OPEN, 10s); on 
timeout, throws and pending
        stays armed at the OLD id, never re-armed [Issue 6]
   COORDINATOR persists CompletedCheckpoint
     -> CompletedCheckpointCodec.encode() called for EVERY job's checkpoint, 
not just dynamic-lookup jobs
        [CheckpointCoordinator.java:1298, verified below] [Issue 5]
     -> CompletedCheckpoint.checkpointIntent field shifts protostuff tag 
numbers for every other field
        declared after it alphabetically [Issue 10]
   ```
   
   **Key findings:**
   - The normal path is the only path. `collect()` is the sole place dimension 
rows are applied and fact rows probed, driven by the ordinary task loop — there 
is no background refresh thread, cache, or executor anywhere in the diff. This 
is a real strength: dimension-apply and fact-probe are strictly sequential on 
one thread and cannot race each other.
   - Where concurrency *does* exist is the gate, not the map: 
`drainGateCommands()` on the task thread (outside the checkpoint lock) versus 
`triggerBarrier -> snapshotGate()` on a Hazelcast operation thread (holding the 
checkpoint lock). The checkpoint lock does not protect the gate because the 
gate mutation lives outside it and `KafkaSourceReader.gateLock` is 
acquired/released across multiple separate critical sections during a single 
snapshot.
   - I independently re-read the primitives, not just the names: 
`checkpointStateLock` in `DynamicLookupFlowLifeCycle` genuinely guards every 
access to 
`barrierAlignments`/`blockedPorts`/`factGateOpened`/`factGateOpening`/`pendingFactGateOpenCheckpointId`
 — correct. `dimensionState` is a plain unlocked `HashMap`, safe only because 
`restoreComplete.complete(null)` establishes a happens-before edge before 
`open()` runs on the task thread — safe by state-machine accident, not by 
declared contract (no threading note exists on the class). 
`KafkaSourceReader.gateOpen` is `volatile` with a real `gateLock` monitor, but 
the *composition* across `snapshotState`/`snapshotGate` is wrong (Issues 2, 3), 
and `SourceReaderBase.splitStates` is a `ConcurrentHashMap`, so the overlap 
yields a torn, weakly-consistent read rather than a 
`ConcurrentModificationException` — which is exactly why this class of bug is 
silent instead of a crash.
   - This is a precise fix for the modelling problem (two-input operator, 
single-threaded state) but the checkpoint-layer and gate changes are 
workarounds that have not converged: both are still unfixed at round 9, at the 
same head.
   
   ## 1.2 Compatibility Impact
   
   **Partially incompatible.**
   
   The feature itself is additive and opt-in: no `dynamic_lookup` block means 
no `DynamicLookupAction`, no gate; `SourceConfig.isDynamicLookupFactGate()` 
defaults false and `KafkaSourceReader.gateOpen` initialises `true`, so an 
ordinary Kafka job is unaffected. The new `PORT_AWARE_LOGICAL_EDGE` serializer 
type id is appended, and `LogicalDag.addEdge`'s new collision check only fires 
for `PortAwareLogicalEdge`, so legacy DAGs skip it. That part is clean.
   
   The checkpoint layer is not scoped to the feature. I confirmed directly: 
`CheckpointCoordinator.completePendingCheckpoint()` calls 
`CompletedCheckpointCodec.encode(completedCheckpoint, serializer)` 
unconditionally for every non-restored checkpoint type, for every job in the 
cluster, not only dynamic-lookup jobs. Decoding falls back to the legacy raw 
serializer when the magic bytes are absent, so *reading* old checkpoints with 
new code works — but a downgrade to a prior release leaves every checkpoint 
written under this release unreadable. Worse, and independently confirmed by 
reading `CompletedCheckpoint.java` and `ProtoStuffSerializer.java`: 
`CompletedCheckpoint` has no `@Tag` annotations on any field, and the 
serializer derives its schema via `RuntimeSchema.createFrom(...)` 
(`serializer-protobuf/.../ProtoStuffSerializer.java:52`), which protostuff 
assigns by alphabetical field name in the absence of explicit tags. Inserting 
`checkpointIntent` between `checkpointId` and `checkp
 ointType` shifts every subsequent field's tag number by one. That means the 
legacy-fallback path in `CompletedCheckpointCodec.decode()` cannot correctly 
deserialize pre-PR checkpoint bytes with the new schema either — restore is 
broken in both directions, for every job, not only dynamic-lookup jobs. None of 
this is recorded in `docs/en/introduction/concepts/incompatible-changes.md`, 
which is not touched by this PR's diff.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   No new thread pools or executors — nothing to leak on that axis. No new 
connection handling to the dimension source beyond the source's own ordinary 
lifecycle.
   
   Memory is nominally bounded but the bound is enforced by a mechanism that 
cannot be exercised at realistic scale: `applyDimension()` 
(`DynamicLookupFlowLifeCycle.java:252-276`) calls 
`enforceDimensionStateBudget()` (`:506-508`), which calls 
`serializeDimensionStatePayload()` (`:494-504`) — a full `ObjectOutputStream` 
walk of the entire `dimensionState` map, on every single dimension row. I read 
this directly; it is exactly as described: ingest cost is O(n²) in both CPU and 
allocation as the map grows, which makes the feature effectively unusable well 
before it reaches its own documented 512 MiB budget. Separately, the 
resident-budget branch in the same method is unreachable, because the parser 
and `DynamicLookupConfig` both already require `resident >= logical`, so any 
payload large enough to trip resident has already thrown on the logical check 
first — a mandatory config option with no observable effect.
   
   ## 1.4 Error Handling and Logging
   
   The failure policy for defined-unsupported conditions (schema change, 
non-INSERT fact rows, dangling `UPDATE_BEFORE`, dimension PK update, budget 
exceeded) is coherent: it fails the job loudly with a specific message. Two 
places break that pattern and both fail silently:
   
   - Dimension-source unavailability has no handling — no timeout, no retry, no 
degraded-mode signal. It cannot deadlock (barriers still flow independently of 
`pollNext`), but if the dimension source stalls before the gate opens, the job 
idles indefinitely with nothing in the logs.
   - INNER-join misses during the warm-up window between gate-open and the 
dimension map becoming complete are silently dropped (`projectFact`, 
`DynamicLookupFlowLifeCycle.java:296-299`) — no metric, no log line, not even 
at DEBUG. In a financial data-integration pipeline this is exactly the kind of 
loss that surfaces during a reconciliation weeks later with no trail back to 
the cause.
   
   **Issue N: New test fixture is invalid HOCON; five of sixteen cases have 
never executed — this is the current CI red at `f243d64d`**
   - **Location:** 
`seatunnel-engine/seatunnel-engine-client/src/test/java/org/apache/seatunnel/engine/client/MultipleTableJobConfigParserTest.java:716,720`
 (`dynamicLookupConfig` helper), `ConfigFactory.parseString` at `:707`.
   - **Problem description:** I read the fixture directly. Lines 716 and 720 
each put two HOCON key/value pairs on the same physical line inside an object 
literal with no separating comma: `schema = { fields { fact_id = "int" 
fact_name = "string" } }`. HOCON needs a newline or a comma between sibling 
fields on one line; this string has neither, so `ConfigFactory.parseString` 
throws `ConfigException.Parse` before the parser under test is ever 
constructed. I re-verified this is still live on a real, fresh CI attempt: fork 
run `31790621314`, `run_attempt=3`, started 2026-08-18T01:42:33Z (two days 
after round 8, same head `f243d64d`) — `unit-test (11, windows-latest)` job 
`95567858216` fails with `Tests run: 16, Failures: 4, Errors: 1`, all traced to 
`MultipleTableJobConfigParserTest.dynamicLookupConfig(MultipleTableJobConfigParserTest.java:707)`,
 `org.apache.seatunnel.shade.com.typesafe.config.ConfigException$Parse`. The 
other three `unit-test` matrix jobs were cancelled by fail-fast,
  not independently failed — this is deterministic (the string is built with 
explicit `\n`, no line-ending sensitivity) and will reproduce on every 
platform. The unrelated `connector-sensorsdata-it` flake round 8 flagged as an 
infra issue is gone on this rerun (71 success / 9 skipped / 3 cancelled / 1 
failure this time), which corroborates that diagnosis.
   - **Potential risk:** All five dynamic-lookup parser tests share this 
helper, so 100% of the parser/validation coverage for this feature — join-key 
arity mismatch, join-key type mismatch, empty `required-capability`, invalid 
byte-size values, and LEFT-join nullability — is currently non-executing. CI is 
red on this PR's own new tests, not on infrastructure.
   - **Best improvement:** Add the missing commas on both lines, then confirm 
the four `assertThrows` cases fail for the intended `JobDefineCheckException` 
reason rather than merely throwing something.
   - **Severity:** High (blocks CI/merge as shipped; does not by itself 
indicate a compile break — the module compiles, only the fixture data is 
malformed).
   
   **Issue N+1: Checkpoint envelope changes the wire format for every job in 
the cluster, with no downgrade path or incompatible-changes entry** — see 1.2 
above for the verified mechanism. **Severity: High.**
   
   **Issue N+2: `checkpointIntent` shifts protostuff tag numbers so legacy 
checkpoints are misread by the new schema, and new checkpoints are unreadable 
by old code** — see 1.2 above; independently re-derived against 
`ProtoStuffSerializer.java:52` (`RuntimeSchema.createFrom`) and the field 
declaration order in `CompletedCheckpoint.java:31-51`, which has no `@Tag` 
annotations anywhere in the class. **Severity: High**, and combined with the 
previous item this is the most consequential defect in the PR for a 
checkpointing engine.
   
   **Issue N+3: `gateOpen` re-read outside `gateLock` after the branch was 
decided under it** — `KafkaSourceReader.java:149-151` (`snapshotState`): the 
ternary at `:144-147` chooses under `synchronized(gateLock)`, but the guard `if 
(!gateOpen)` at `:150` reads the field again after the lock is released, and 
`activateStagedSplits()` on the task thread can flip it in between. Verified 
directly. Not a regression for ungated jobs (gate permanently open there), but 
live for every gated job — can register offsets for splits that were never 
actually read, or skip registering offsets that were. **Severity: High.**
   
   **Issue N+4: `snapshotGate()` composes a three-field state (`sourceSplits`, 
`gateOpen`, `stagedNoMoreSplits`) from two separate lock acquisitions** — 
`KafkaSourceReader.java:210-222`: `snapshotState(checkpointId)` 
acquires+releases `gateLock` once internally, then `snapshotGate` re-acquires 
it separately to read `gateOpen`/`stagedNoMoreSplits`. Verified directly 
against the current source. `activateStagedSplits()` can interleave between the 
two acquisitions and flips `stagedNoMoreSplits=false` as it drains — so a 
restored reader can silently lose a "no more splits" signal it had already 
received, and a bounded fact source then never signals completion after 
recovery: this corrupts recovery state, it does not merely degrade throughput. 
**Severity: High.**
   
   **Issue N+5: Entire dimension map re-serialised on every dimension row** — 
see 1.3; verified directly at `DynamicLookupFlowLifeCycle.java:252-276` (call 
site) and `:494-524` (the O(n) `ObjectOutputStream` walk invoked on every 
mutation). **Severity: High.**
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   Formatting, license headers, import hygiene are fine; no wildcard imports, 
no `System.out.println`, no single-line inline Javadoc. Documentation on new 
core methods/fields is good where it exists (`checkpointStateLock`, 
`MAX_RECORDS_PER_PORT_DRAIN`, `factGateOpening`, 
`serializeDimensionStatePayload` — all explain *why*). The one systemic gap: 
nothing documents which methods run on the task thread, which on a Hazelcast 
operation thread, and which on the checkpoint-callback thread, even though 
`dimensionState`'s safety depends entirely on that ordering. **Medium** — 
recommend adding a class-level threading note to `DynamicLookupFlowLifeCycle` 
and `FactSourceGateCapability` before merge, since it is exactly the 
documentation a future maintainer needs to avoid reintroducing Issues 2/3's 
class of bug elsewhere.
   
   ## 2.2 Test Coverage and Test Stability
   
   Nine new test files, 26 test methods. 
`LogicalEdgeSerializationCompatibilityTest` (byte-exact golden vectors plus an 
unknown-format-version rejection case) and `CompletedCheckpointCodecTest`'s 
legacy-payload round-trip are the right shape for pinning a wire contract. 
`ExecutionPlanGeneratorTest`'s rejection cases genuinely lock down 
DAG-validation rules.
   
   Against the runtime itself, coverage has real gaps I confirmed by reading 
the test files: `collect()`, `processRecord()`, `applyDimension()`, 
`projectFact()`, `alignBarrier()`, `restoreState()` and both budget overloads 
are reached by zero tests — the two control-plane tests poke private fields via 
reflection and assert the same fields afterward, without driving a record, a 
barrier, or a real queue. `applyGateCommand` is never invoked by any test 
(including `ABORT` and the `CLOSE`-after-activation `IllegalStateException` 
branches). There is no `seatunnel-e2e` fixture for dynamic lookup at all. A 
repo-wide check for `Thread(`, `ExecutorService`, `CountDownLatch`, 
`CyclicBarrier`, `Awaitility`, `@RepeatedTest` across the new test files 
returns nothing — every one of the 26 tests is single-threaded, which is 
exactly why Issues 2 and 3 (pure interleaving defects between the task thread 
and a Hazelcast operation thread) cannot be caught by this suite regardless of 
how many rounds i
 t goes through.
   
   Flaky-pattern check on the new tests: no `Thread.sleep`/unbounded waits 
found — good. Two lower-severity notes: `PendingCheckpointFinalizeTest` bounds 
a synchronously-completed future with a 1s wall-clock `get(...)` where 
`getNow()` would suffice; `CheckpointPlanTest` indexes `HashMap`-backed 
structures with `get(0)`/`keySet().iterator().next()`, which only passes 
because each collection happens to be size-1 today.
   
   **Rating: High risk.** This is a formal issue, not just a note, because it 
is the direct explanation for why Issues 2/3 (High-severity recovery 
correctness bugs) currently ship with zero test coverage that could ever detect 
them, and because the suite as it stands is currently red (Issue on HOCON 
fixture above) rather than merely thin.
   
   **Live CI status at `f243d64d` (re-verified 2026-08-19, independent of round 
8's 2026-08-16 read):** apache-side `Build` check is a pointer to the fork run. 
Fork run `31790621314` (`run_attempt=3`, 
`run_started_at=2026-08-18T01:42:33Z`): 71 success, 9 skipped, 3 cancelled, 1 
failure. The failure is `unit-test (11, windows-latest)` (job `95567858216`) — 
I fetched and grepped the raw log directly rather than trusting the summary; it 
fails at the surefire phase of `seatunnel-engine-client` with `Tests run: 16, 
Failures: 4, Errors: 1` in `MultipleTableJobConfigParserTest`, all four 
`AssertionFailedError: Unexpected exception type thrown ... but was 
ConfigException.Parse`, tracing to the same malformed fixture. No compilation 
error anywhere in the log — the code itself compiles cleanly at this head. The 
three other `unit-test` matrix jobs were cancelled by fail-fast, so this is not 
platform-specific; it will reproduce everywhere the suite runs. This is the 
PR's own new test fixture
  failing, not CI flakiness and not a Zeta compile break.
   
   ## 2.3 Documentation Updates
   
   Both `docs/en` and `docs/zh` pages exist and are registered in 
`docs/sidebars.js`. Content issues carried from round 8 remain unfixed at this 
head (I did not re-derive these from scratch this round, since the doc files 
are unchanged from what round 8 already quoted line-for-line): the worked 
config example uses a bare `table = "customers"` where the parser requires a 
fully-qualified `mydb.customers`-style name; the nullability section states the 
opposite of what `MultipleTableJobConfigParser.java`'s `nullable = ... joinType 
== LEFT || sourceColumn.isNullable()` logic actually does; and there is no 
mention of the staleness/warm-up window, the silent INNER-join drop, or that 
`uid` (undocumented) is the checkpoint-identity anchor across restarts.
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance
   
   Modelling this as a first-class two-input `DynamicLookupAction` instead of a 
Transform is the right call — a two-input operator does not fit the transform 
contract. Keeping dimension state single-threaded on the operator's own thread 
removes an entire class of concurrency risk by construction. Fencing M0 down to 
`IN_MEMORY`/`ttl=NONE`/`APPEND_ONLY`/`FAIL` is exactly how a risky engine 
feature should land. **Classification: precise fix for the modelling problem.** 
Where elegance slips is the gate: its state is split across three owners 
(operator's `factGateOpened`, reader's `gateOpen`/`activateOnOpen`, and a 
one-slot command queue) with no single authority — the direct cause of Issues 
2, 3 and 6. That part reads as a **temporary workaround** that has not 
converged across nine review rounds.
   
   ## 3.2 Maintainability
   
   Roughly 750 lines of channel-layer scaffolding (`ChannelEnvelope`, 
`ChannelAttemptId`, `ChannelEnvelopeDeduplicator`) have no production consumer 
— referenced only by their own tests, confirmed by grep. 
`ChannelEnvelopeDeduplicator.appliedDigests` is an unbounded 
`ConcurrentHashMap` with `putIfAbsent`-only growth; it is a latent rather than 
live leak today, but it is a loaded gun for whoever wires it up next with no 
warning in the class that it needs a bound.
   
   ## 3.3 Extensibility
   
   The extension seams are well chosen — capability interfaces on the reader 
side, versioned state envelopes with magic+digest, an operator-scoped 
checkpoint identity already threaded through the coordinator. Moving from M0's 
in-memory backend to disk-backed state, or from same-task-group queues to 
cross-worker channels, should not require re-cutting these interfaces.
   
   ## 3.4 Historical-Version Compatibility
   
   This is the section that blocks the merge, and it is unchanged since round 8 
at the same head: restore is broken in both directions (Issues N+1/N+2 above), 
the checkpoint change is release-wide rather than feature-scoped, and 
`incompatible-changes.md` is untouched. For a checkpointing engine this is the 
highest-consequence defect category available, ahead of the concurrency 
findings.
   
   # 4. Issue Summary
   
   Numbering continues from round 8's numbering (in brackets) for author 
continuity across rounds; "round 9" marks what I re-verified directly against 
the current head and/or fresh CI rather than carrying forward unread.
   
   | # | Issue | Location | Severity | Status |
   | --- | --- | --- | --- | --- |
   | 1 | Invalid HOCON in test fixture; 5/16 cases never execute; current CI 
red | `MultipleTableJobConfigParserTest.java:707,716,720` | High | Re-verified 
round 9 (fresh CI log, fresh fixture read) |
   | 2 [30] | `gateOpen` re-read outside `gateLock` after branch decided under 
it | `KafkaSourceReader.java:144-151` | High | Re-verified round 9 (direct 
source read) |
   | 3 [34] | `snapshotGate()` composes splits+flags from 2 separate lock 
acquisitions | `KafkaSourceReader.java:210-222` | High | Re-verified round 9 
(direct source read) |
   | 4 [3] | Entire dimension map re-serialised on every dimension row — O(n^2) 
| `DynamicLookupFlowLifeCycle.java:252-276,494-524` | High | Re-verified round 
9 (direct source read) |
   | 5 [31] | Checkpoint envelope applied to every job; no downgrade path/doc | 
`CheckpointCoordinator.java:1298` | High | Re-verified round 9 (direct source 
read) |
   | 10 [32] | `checkpointIntent` shifts protostuff tags; legacy checkpoints 
misread | `CompletedCheckpoint.java:31-51`, `ProtoStuffSerializer.java:52` | 
High | Re-verified round 9 (traced to `RuntimeSchema.createFrom`) |
   | 6 [28] | Failed gate open never retried; permanent silent stall | 
`DynamicLookupFlowLifeCycle.java:374-447` | Medium | Carried from round 8, 
unread this round |
   | 7 | Resident-budget check unreachable dead code | 
`DynamicLookupFlowLifeCycle.java:518-524` | Medium | Re-verified round 9 (read 
alongside Issue 4) |
   | 8 | INNER-join misses silently dropped, no counter/log | 
`DynamicLookupFlowLifeCycle.java:296-299` | Medium | Re-verified round 9 
(direct source read) |
   | 9 [11] | `CheckpointIntent` write-only; SHA-256 over all state per 
checkpoint | `PendingCheckpoint.java:183-217` | Medium | Carried from round 8, 
unread this round |
   | 11 | `stateSize` still counts elements, not bytes | 
`PendingCheckpoint.java:149` | Medium | Carried from round 8, unread this round 
|
   | 12 | Doc example `dimension.table = "customers"` cannot run | 
`docs/en,zh/.../dynamic-lookup.md` | Medium | Carried from round 8, unread this 
round |
   | 13 | Doc nullability claim states the opposite of the code | 
`docs/en,zh/.../dynamic-lookup.md` | Medium | Carried from round 8, unread this 
round |
   | 14 [12] | ~750 lines unwired scaffolding; deduplicator map unbounded | 
`ChannelEnvelopeDeduplicator.java:36` | Medium | Re-verified round 9 (grep for 
consumers) |
   | 15 | No concurrent/restart-mid-load/source-failure/E2E test | new test 
files | Medium | Re-verified round 9 (grep for `Thread(`/`Awaitility`/etc.) |
   | 16 [35] | `collect()` idles at flat 100ms; 10ms poll per record | 
`DynamicLookupFlowLifeCycle.java:129-131` | Low | Carried from round 8, unread 
this round |
   | 17 [33] | Native `ObjectInputStream` on checkpoint bytes, no class filter 
| `DynamicLookupFlowLifeCycle.java:557` | Low | Carried from round 8, unread 
this round |
   | 18 | Undocumented `uid` silently controls checkpoint identity | 
`MultipleTableJobConfigParser.java:700` | Low | Carried from round 8, unread 
this round |
   | 19 | Threading contract undocumented on operator/capability interface | 
`DynamicLookupFlowLifeCycle.java` | Low | Re-verified round 9 (see 2.1) |
   
   # 5. Merge Recommendation
   
   ### Conclusion: Not recommended for merge
   
   **Blockers (must be fixed, all re-verified against the current head or fresh 
CI this round):**
   1. Issues 5 and 10 — checkpoint restore is broken in both directions and the 
change reaches outside this feature's blast radius to every job in the cluster; 
no `incompatible-changes.md` entry exists. For a checkpointing engine this is 
the single most serious defect and it alone should block merge.
   2. Issue 1 — the PR's own new tests do not run; CI is red on this PR's own 
code, confirmed on a fresh rerun (2026-08-18) two days after round 8, at the 
unchanged head.
   3. Issues 2 and 3 — the Kafka gate's snapshot path reads composite state 
non-atomically, which corrupts recovery state (a bounded job can hang forever 
after restore) rather than merely losing throughput.
   4. Issue 4 — per-row full re-serialisation makes the feature unusable at any 
dimension size that matters in practice, well before its own documented budget.
   
   **Recommended before merge (non-blocking but should not wait for a 
follow-up):** Issues 6, 7, 8, 9, 11, 14, 15 — of these, Issue 8 (silent 
INNER-join data loss with no metric, in a financial-grade data-integration 
engine) and Issue 15 (zero concurrent/E2E coverage, which is why Issues 2/3 
have survived nine rounds) deserve the most weight.
   
   **Recommended follow-up:** Issues 12, 13, 16, 17, 18, 19.
   
   **Overall assessment.** The core modelling decisions remain sound and I want 
that on record again: a first-class two-input action instead of a Transform, 
single-threaded dimension state, and an aggressively fenced M0 feature surface 
are all the right calls, and the capability-negotiation and versioned-envelope 
seams give the follow-on work (disk-backed state, cross-worker channels) clean 
places to attach. What is not ready, unchanged across two additional CI 
attempts and one more review round since 2026-08-16, is the checkpoint layer 
and the gate. Neither has a test that could catch its own defect class — the 
checkpoint tests round-trip within a single schema version, and every gate test 
is single-threaded. My read, independent of the eight prior rounds' framing but 
arriving at the same place from a direct re-check of the source: the 
lookup-operator logic has converged, the checkpoint-envelope and 
`CheckpointIntent` changes underneath it have not, and shipping them in one PR 
mea
 ns a checkpoint-compatibility blocker is gating an otherwise-landable feature. 
I would split the checkpoint-envelope/`CheckpointIntent` work into its own 
change with dedicated cross-version restore tests, and let this PR land the 
lookup operator and gate once Issues 1-4 are fixed on their own merits.
   
   This posts as an issue comment rather than a pull-request review: it is my 
own PR, and GitHub does not permit self-approval on it.
   


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