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]