DanielLeens commented on PR #11559:
URL: https://github.com/apache/seatunnel/pull/11559#issuecomment-5379409363
# What Problem Does This PR Solve?
- **User pain point**: Apache SeaTunnel's Zeta engine has no engine-native
way to join an append-only "fact" stream (e.g., Kafka) against a CDC
"dimension" stream inside a single job while keeping fact-position ownership,
checkpoint/restore identity, source ordering, and dimension-state memory
bounded. Today this requires bolting together separate jobs or an external
lookup service.
- **Fix approach**: A first-scope ("M0") `dynamic_lookup { ... }` runtime: a
two-input `DynamicLookupAction` (fact input + dimension input), an in-memory
keyed dimension map fed by a CDC changelog, a Kafka fact-source "gate" that
holds the fact source closed until the dimension side has a head start, and a
versioned/digest-checked checkpoint envelope scoped to this feature only so
ordinary jobs are byte-for-byte unaffected.
- **One-sentence summary**: Adds an opt-in, in-memory dynamic-lookup join
operator to the Zeta engine with its own checkpoint/gate/state-budget machinery
that is isolated from the ordinary job checkpoint format.
This is round 13 on this PR. I read the full diff from scratch against the
current head and independently re-verified the two substantive open items from
round 12 (a disputed security claim, and CI status) from source rather than
trusting either round's prose, per the review protocol.
# 1. Code Change Review
## 1.1 Core Logic Analysis
**Head check**: current head `f40ad4eaa1e1930ca9fa29003f632819345837e1` is
unchanged since round 10 (2026-08-20T10:40 UTC) — confirmed via `headRefOid`.
No new commits since then; the only new activity between round 12 and now is
nothing (round 12 is still the latest activity). This round re-verifies round
12's conclusions from source rather than re-stating them.
**1) Re-verifying the disputed deserialization-safety claim (SEZ9 vs. round
12).**
SEZ9's 2026-08-21T06:54 review claims gate-state restore in
`SourceFlowLifeCycle` and dimension-state restore in
`DynamicLookupFlowLifeCycle` "still call `readObject()` with no class
restrictions," citing `DynamicLookupFlowLifeCycle.java:558/566` and
`SourceFlowLifeCycle.java:796`. Round 12 disputed this, citing
`DynamicLookupFlowLifeCycle.java:641-642` as the actual `readObject()` call
sites, guarded by a class-allowlisting `ObjectInputStream` subclass.
I fetched both files at the current head directly (not from either review's
prose) and read the full method bodies:
```java
// DynamicLookupFlowLifeCycle.java, restoreDimensionState() — verified at
current head
try (ObjectInputStream inputStream =
new DynamicLookupStateObjectInputStream(new
ByteArrayInputStream(payload))) {
int entryCount = inputStream.readInt();
...
for (int i = 0; i < entryCount; i++) {
LookupKey key = (LookupKey) inputStream.readObject(); // line
641
SeaTunnelRow row = (SeaTunnelRow) inputStream.readObject(); // line
642
...
}
}
private static final class DynamicLookupStateObjectInputStream extends
ObjectInputStream {
@Override
protected Class<?> resolveClass(ObjectStreamClass descriptor) throws
IOException, ClassNotFoundException {
String className = descriptor.getName();
if (isAllowedClass(className)) {
return super.resolveClass(descriptor);
}
throw new IOException("Rejected dynamic lookup dimension state
class: " + className);
}
private static boolean isAllowedClass(String className) {
// explicit allowlist: java.lang.*, java.math.*, java.sql.*,
java.time.*, java.util.*,
// org.apache.seatunnel.api.table.type.*, LookupKey,
ByteArrayKeyPart, primitive arrays
}
}
```
`SourceFlowLifeCycle.java` has the exact same pattern
(`SourceGateObjectInputStream` overriding `resolveClass` before the
`readObject()` call). Both classes reject any class not on the explicit
allowlist by throwing `IOException` — this is a real `resolveClass`-based
allowlist, not a bolt-on filter or a no-op. Both restore paths also verify a
SHA-256 digest over the raw payload bytes *before* deserialization even begins
(`DynamicLookupStateEnvelope`, `CompletedCheckpointCodec`), so a bit-flipped
payload is rejected before `readObject()` runs at all.
**My independent conclusion: SEZ9's claim does not hold at the current head
`f40ad4eaa1e1`.** The line numbers SEZ9 cited (558/566/796) do not correspond
to `readObject()` calls in the current files (I confirmed the actual
`readObject()` sites are at 641-642 and ~798 respectively) — this is consistent
with round 12's theory that SEZ9 was looking at a pre-round-10 revision, since
the commit at this exact head (`f40ad4eaa`) is titled "Fix dynamic lookup
checkpoint and gate recovery" and round 10's own history describes adding
exactly this class-allowlist hardening. @SEZ9 — could you confirm which commit
SHA you had checked out when you found the unguarded `readObject()` calls? If
it was an earlier revision, this should already be resolved at `f40ad4eaa1e1`;
if you're still seeing an unguarded path at this exact SHA, please point me at
the exact file/line so I can re-check — I want to make sure this genuinely
High-severity class of bug doesn't slip through if I've misread somet
hing.
One real, non-blocking observation on the same code: `isAllowedClass`
allowlists whole `java.util.*`/`java.time.*`/`java.sql.*`/`java.math.*`
packages rather than the specific concrete classes
`serializeDimensionStateEntry` actually writes (e.g. `ArrayList`, `HashMap`,
`BigDecimal`, `LocalDateTime`). This is far safer than unrestricted
deserialization, but broader than the minimum needed — see Issue 2 below.
**2) Re-verifying checkpoint format compatibility for ordinary
(non-dynamic-lookup) jobs.**
This is the single highest-stakes correctness claim in this PR — if wrong,
it would corrupt checkpoint restore for every existing SeaTunnel job, not just
dynamic-lookup ones. I read `CompletedCheckpointCodec.java` directly:
```java
public static byte[] encode(CompletedCheckpoint checkpoint, Serializer
serializer) throws IOException {
byte[] payload = serializer.serialize(checkpoint);
if (checkpoint.getCheckpointIntent().isNormalCheckpoint()) {
return payload; // <-- byte-for-byte identical to
pre-PR output
}
... // magic + version + length + payload + SHA-256, only for
dynamic-lookup checkpoints
}
public static CompletedCheckpoint decode(byte[] states, Serializer
serializer) throws IOException {
int magic = buffer.getInt();
if (magic != MAGIC) {
return serializer.deserialize(states, CompletedCheckpoint.class);
// legacy fallback
}
... // versioned envelope path
}
```
Confirmed: `encode()` returns the untouched legacy `payload` for any
checkpoint whose `CheckpointIntent.isNormalCheckpoint()` is true, and
`CheckpointIntent.normal(...)` is exactly what
`PendingCheckpoint.createCheckpointIntent()` produces whenever no action state
contains a dynamic-lookup envelope magic. `decode()` only takes the
new-envelope path when the first 4 bytes match the specific `MAGIC` constant
(`0x53544350`); otherwise it falls straight through to the legacy deserializer,
so already-written checkpoints from before this PR remain readable. I also
checked `@Tag` numbering on `CompletedCheckpoint`'s fields is unchanged from
pre-PR (new fields appended after existing tags, not renumbered) — this is what
makes the "legacy raw payload" path actually round-trip correctly.
One low-probability theoretical note I want to flag for completeness rather
than as an actionable issue: `decode()`'s dispatch is a 4-byte magic sniff on
the legacy path with no independent version marker — if the legacy
`Serializer`'s output for some future `CompletedCheckpoint` shape ever happened
to start with the exact bytes `0x53544350`, it would be misrouted into the
versioned-envelope path and fail the format-version or digest check. Given the
legacy serializer is Hazelcast's `IdentifiedDataSerializable` binary format
(factory-id/class-id integers, not free-form bytes), this is astronomically
unlikely in practice and I'm not raising it as a formal issue — flagging only
so it's a documented, consciously-accepted trade-off rather than an unnoticed
one.
**3) CI status — independently extended past round 12's "assumption, not
confirmed."**
Round 12 confirmed the Windows `unit-test` failure is a known,
already-fixed-elsewhere flake (`MultiTableSinkWriterSchemaChangeBroadcastTest`,
not touched by this PR, already fixed on `dev` at `22da84a418` / PR #11726, and
this branch is 103 commits behind `dev` so it predates that fix), but
explicitly flagged the other three failing jobs (`rocketmq-connector-it`,
`paimon-connector-it`, `all-connectors-it-5`) as an unverified assumption. I
individually pulled and read each of the three remaining failing jobs' logs
from the fork run (`DanielLeens/seatunnel` run `32328182757`):
- `rocketmq-connector-it`: `RocketMqIT.testRocketMqTimestampToConsole` —
`ConditionTimeout` after a `RemotingSendRequestException` in the Flink
taskmanager container talking to the RocketMq broker container — a
container-networking timeout, unrelated to this diff.
- `paimon-connector-it`: `PaimonSinkCDCIT.testSinkWithIncompatibleSchema` —
`ConditionTimeout` on an Awaitility assertion — a Paimon schema-evolution IT,
unrelated to this diff.
- `all-connectors-it-5`: `DatabendCDCSinkIT` — `ContainerLaunchException`
(Testcontainers startup failure) — a Databend CDC sink IT, unrelated to this
diff.
None of `connector-rocketmq`, `connector-paimon`, or `connector-databend`
appear anywhere in this PR's changed-file list (`git diff dev...HEAD
--name-status` — confirmed via the actual diff, zero matches for any of the
three module paths). Combined with round 12's confirmed Windows finding, **all
four fork CI failures at this head are now confirmed (not assumed) to be
pre-existing/environmental failures unrelated to this PR's own code**,
consistent with the branch being 103 commits stale.
**Runtime path (checkpoint + restore), re-verified at this head, unchanged
from round 10/12:**
```
CHECKPOINT (Hazelcast operation thread)
-> PendingCheckpoint.acknowledgeTask() ->
finalizeState.compareAndSet(ACTIVE, COMPLETED)
(prevents a double-finalize race against a concurrent
abortCheckpoint())
-> toCompletedCheckpoint() -> createCheckpointIntent()
containsDynamicLookupState() ? dynamicLookupFactPositionAnchor(...,
digestActionStates())
: CheckpointIntent.normal(...)
-> CompletedCheckpointCodec.encode(): isNormalCheckpoint() ? raw legacy
bytes : magic+digest envelope
RESTORE
-> CompletedCheckpointCodec.decode(): magic mismatch -> legacy
deserializer (old checkpoints still readable)
-> DynamicLookupFlowLifeCycle / SourceFlowLifeCycle restore
SHA-256 digest check on the envelope payload, then
class-allowlisted ObjectInputStream (verified above) rejects any
non-allowlisted class
```
## 1.2 Compatibility Impact
**Partially incompatible, and correctly scoped.** Ordinary
(non-`dynamic_lookup`) jobs are byte-for-byte unaffected in both write and read
directions, verified directly in 1.1 above (not just asserted). The
incompatibility is intrinsically scoped to jobs that opt into the new
`dynamic_lookup` feature, which did not exist before this PR — there is no
existing job whose behavior this changes.
`docs/en/introduction/concepts/incompatible-changes.md` is still not updated to
describe the dynamic-lookup checkpoint envelope's own forward/backward
semantics (e.g., that a `dynamic_lookup` job's checkpoint cannot be restored by
a pre-PR SeaTunnel version) — recommended before merge, not release-blocking
given the narrow blast radius (opt-in, brand-new feature).
## 1.3 Performance / Side-Effect Analysis
- Dimension-state ingest tracks a state-size budget incrementally; full
serialization only happens at checkpoint/restore time, not per-record.
- `ChannelEnvelopeDeduplicator` is bounded at a fixed capacity (I confirmed
this in the diff), not an unbounded cache.
- No new thread pools, executors, or connections introduced.
- `digestActionStates()` in `PendingCheckpoint.java` (verified at
current-head lines 217-, called from `createCheckpointIntent()` at line 187)
computes a full SHA-256 digest over every action's coordinator + subtask state
bytes on **every** dynamic-lookup checkpoint. I checked whether
`CheckpointIntent.getAnchoredPositionDigest()` (the only reader of this digest)
is called anywhere in production code — grepping the full diff, the **only**
call site is in `PendingCheckpointFinalizeTest.java` (a test assertion). There
is no production consumer that reads this digest back to verify or anchor
anything. See Issue 1 below — this re-confirms round 10/12's finding rather
than just repeating it.
- The idle-poll loop in
`DynamicLookupFlowLifeCycle`/`DynamicLookupCoordinatorTask` uses a flat
`Thread.sleep(IDLE_SLEEP_MS)` with `IDLE_SLEEP_MS = 100L` (confirmed present at
multiple call sites in the diff) rather than an event-driven wakeup or adaptive
backoff — see Issue 2.
## 1.4 Error Handling and Logging
### Issue 1: `CheckpointIntent.anchoredPositionDigest` is computed on every
checkpoint but has no production reader (Medium, carried over — independently
re-confirmed this round)
- **Location**: `PendingCheckpoint.java` (`digestActionStates()`,
`createCheckpointIntent()`); `CheckpointIntent.getAnchoredPositionDigest()`.
- **Problem**: A full SHA-256 digest over all action/subtask state bytes is
computed on every dynamic-lookup checkpoint, but I confirmed (by grepping every
reference to `getAnchoredPositionDigest()` in the diff) the only caller is a
test assertion — no production code path reads this value back to verify or
anchor a fact position, despite the class/field naming ("fact position anchor")
strongly implying it should.
- **Potential risk**: Pure wasted CPU today (proportional to checkpoint
state size, on the checkpoint-coordinator's critical path); more importantly, a
half-wired verification mechanism reads as "this get checked" to a future
maintainer when it currently doesn't, which is a worse trap than not having the
field at all.
- **Best improvement**: Either wire a real consumer (e.g., verify the
anchored digest matches recomputed state on restore, to catch state corruption)
or remove the field and its computation until a consumer exists, per the
M0/first-scope nature of this PR.
- **Severity**: Medium.
- **Raised by another reviewer**: No — carried over from my own round 10/12
findings; independently re-verified this round (I did not simply trust the
prior round's citation).
### Issue 2: Flat 100ms idle-poll sleep, no adaptive backoff (Low, carried
over — independently re-confirmed)
- **Location**: `DynamicLookupFlowLifeCycle.java` /
`DynamicLookupCoordinatorTask.java`, `IDLE_SLEEP_MS = 100L`, multiple
`Thread.sleep(100)` call sites (confirmed present in the diff at the current
head).
- **Problem**: A fixed 100ms poll interval regardless of load or idle
duration.
- **Potential risk**: Minor CPU/latency inefficiency under sustained idle;
not a correctness issue.
- **Best improvement**: Adaptive/capped backoff instead of a flat constant.
- **Severity**: Low.
- **Raised by another reviewer**: No.
### Issue 3: `isAllowedClass` deserialization allowlist is broader than the
concrete types actually written (Low, carried over — independently re-confirmed)
- **Location**: `DynamicLookupFlowLifeCycle.java`
(`DynamicLookupStateObjectInputStream.isAllowedClass`),
`SourceFlowLifeCycle.java` (equivalent class).
- **Problem**: Whole-package allowlisting (`java.util.*`, `java.time.*`,
`java.sql.*`, `java.math.*`) rather than the specific concrete classes written
by the corresponding serialize methods.
- **Potential risk**: Marginally larger attack surface than the minimum
necessary for a tampered checkpoint payload that already passed the SHA-256
digest check (i.e., an attacker with checkpoint-storage write access) — still
far safer than unrestricted deserialization.
- **Best improvement**: Narrow to the exact concrete classes written.
- **Severity**: Low.
- **Raised by another reviewer**: Related to SEZ9's finding, but factually
distinct — SEZ9 claimed these sites are entirely unrestricted, which I found to
be incorrect at the current head (see 1.1); this is a narrower, non-blocking
hardening suggestion on top of a mechanism that already works correctly.
# 2. Code Quality Assessment
## 2.1 Coding Standards
No wildcard imports, no `System.out.println`, ASF headers present on all new
files, standard multi-line Javadoc (no single-line `/** ... */`). The new
`*ObjectInputStream` inner classes each carry a purpose comment.
`PendingCheckpoint`'s new `FinalizeState` enum + `AtomicReference` (guarding
against a completion/abort race between `acknowledgeTask()` and
`abortCheckpoint()`) is a good, minimal concurrency-safety addition I hadn't
seen called out explicitly in prior rounds — it's a `compareAndSet`-gated state
machine with exactly two committing transitions (`ACTIVE`→`COMPLETED`,
`ACTIVE`→`ABORTED`), which is the right shape for this kind of "exactly one of
two terminal outcomes" race.
## 2.2 Test Coverage and Test Stability
**Rating: Risk present** (unchanged from round 10/12, for the same reason).
The new unit tests themselves (`DynamicLookupFlowLifeCycleControlPlaneTest`,
`DynamicLookupFlowLifeCycleDataPlaneTest`, `ChannelEnvelopeDeduplicatorTest`,
`LogicalEdgeSerializationCompatibilityTest`, `PendingCheckpointFinalizeTest`,
`CompletedCheckpointCodecTest`, `ExecutionPlanGeneratorTest` additions,
`KafkaSourceReaderGateTest`) are deterministic — I specifically checked for
`Thread.sleep`/`Awaitility` usage inside these test files and found none; they
exercise the poll/tick logic synchronously rather than waiting on real timers,
and `LogicalEdgeSerializationCompatibilityTest` uses golden-byte-vector
fixtures rather than round-trip-only assertions, which is a good pattern for a
wire-format compatibility test. The residual risk is external to the test code
itself: no dedicated E2E test exercises a real `dynamic_lookup` job end-to-end
(submit, checkpoint, restart mid-load, verify no duplicate/lost dimen
sion state) — I confirmed only one E2E-adjacent file is touched by this PR
(`MongodbCDCIT.java`, an incidental interface-signature adaptation, not a
dynamic-lookup feature test) and no new `*IT.java` under `seatunnel-e2e` was
added for this feature. For a checkpoint/restore-touching engine feature, a
restart-mid-load E2E is the kind of test that catches interaction bugs unit
tests structurally cannot.
## 2.3 Documentation Updates
`docs/en/architecture/features/dynamic-lookup.md` and the `zh` equivalent
are present and current for the feature's own configuration/behavior.
`docs/en|zh/introduction/concepts/incompatible-changes.md` still has no entry
for the dynamic-lookup checkpoint envelope's cross-version restore semantics
(see 1.2) — recommended before merge.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
**Precise fix for the checkpoint-format-scoping problem specifically**:
gating the new envelope behind `CheckpointIntent.isNormalCheckpoint()` and
preserving `@Tag` ordering is the structurally right way to add a
feature-scoped persistence format without touching the existing one. The
class-allowlist `ObjectInputStream` pattern for restore-path deserialization is
also the right primitive (verified working, not just claimed, in 1.1). This is
still an M0/first-scope feature by its own PR description, and the open items
(Issues 1-3, plus the missing E2E) are consistent with that scope rather than
being missed fundamentals.
## 3.2 Maintainability
`ChannelEnvelope`/`ChannelEnvelopeDeduplicator` remain unwired to any
production consumer — I confirmed this again this round by searching the full
diff for references outside their own definition and test files; there are
none. This isn't a live leak or correctness risk since it's inert, but it reads
as either future work that should be tracked explicitly or dead code that
should be deferred to the PR that actually wires it up.
## 3.3 Extensibility
`PortAwareAction`/`PortAwareLogicalEdge`/`ExchangeDescriptor` (the
two-input-port DAG plumbing) look reusable for future multi-input operators
beyond dynamic lookup specifically, which is a reasonable investment for a
first engine-level multi-input feature.
## 3.4 Historical-Version Compatibility
Ordinary-job checkpoint restore is fully compatible in both directions,
verified directly in 1.1 (not merely asserted). `dynamic_lookup`-job restore
across SeaTunnel versions is an intrinsic one-way door for a brand-new opt-in
feature (no prior version could have written this state), not a regression
against any existing release.
# 4. Issue Summary
| # | Issue | Location | Severity |
| --- | --- | --- | --- |
| — | SEZ9's "unrestricted `readObject()`" claim |
`DynamicLookupFlowLifeCycle.java`, `SourceFlowLifeCycle.java` | — Independently
re-verified as **not reproducing at the current head**; both sites use a
`resolveClass`-allowlisting `ObjectInputStream` subclass plus a
pre-deserialization SHA-256 digest check. Awaiting SEZ9's confirmation of which
commit they reviewed. |
| — | Branch 103 commits behind `dev`; all 4 fork CI failures confirmed
pre-existing/environmental | branch state; `unit-test(windows)`,
`rocketmq-connector-it`, `paimon-connector-it`, `all-connectors-it-5` | High
(blocks a trustworthy CI signal, not a defect in this diff) |
| 1 | `anchoredPositionDigest` computed every checkpoint, no production
reader | `PendingCheckpoint.java:187-217`, `CheckpointIntent.java` | Medium |
| 2 | Flat 100ms idle-poll sleep | `DynamicLookupFlowLifeCycle.java` /
`DynamicLookupCoordinatorTask.java` | Low |
| 3 | Deserialization allowlist broader than necessary |
`DynamicLookupFlowLifeCycle.java`, `SourceFlowLifeCycle.java` | Low |
| — | No restart-mid-load / concurrency-stress E2E for the new feature |
test files | Medium |
| — | `incompatible-changes.md` not updated for the dynamic-lookup
checkpoint envelope | docs | Medium |
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
**1. Blockers — must be fixed before merge:**
- Rebase/sync this branch onto current `apache/seatunnel:dev` (103 commits
behind) and get a fresh CI run on the rebased head. I independently confirmed
all 4 current fork CI failures are pre-existing/environmental and unrelated to
this PR's own changed files (see 1.1) — but a trustworthy green (or
triaged-red-for-known-reasons) signal on a current-`dev`-based head is still
required before merge, since a stale base means CI isn't actually validating
this diff against what it will land on.
- @SEZ9's High-severity deserialization-safety finding needs to be resolved
as a discussion, not just left as a disagreement: either confirm the commit you
reviewed was stale (in which case this is already fixed at `f40ad4eaa1e1`, as I
independently verified), or point at a specific line at the current head that's
still unguarded so it can be fixed before merge. Given how serious an
unrestricted-deserialization gadget-chain entry point would be if it were real,
this shouldn't merge while there's an open, unreconciled dispute about it —
even though my own independent read of the source says it's already fixed.
**2. Recommended fixes — non-blocking (this PR is explicitly
M0/first-scope):**
- Issue 1 — wire a real consumer for `anchoredPositionDigest` or remove the
computation until one exists.
- Issue 2 — adaptive backoff instead of the flat 100ms poll.
- Issue 3 — narrow the deserialization allowlist to concrete types actually
written.
- Add the `incompatible-changes.md` entry for the dynamic-lookup checkpoint
envelope.
- At least one restart-mid-load E2E test before this feature graduates past
M0/draft.
**Overall assessment.** This is a large, genuinely complex engine feature
(new DAG port/edge model, feature-scoped checkpoint envelope, gated source,
bounded dimension-state map) and, on this round's independent re-verification,
its two highest-stakes correctness properties both hold up under direct source
inspection rather than just being asserted: ordinary-job checkpoint
compatibility is real (traced through `CompletedCheckpointCodec`), and the
restore-path deserialization hardening SEZ9 disputed is also real (traced
through both `ObjectInputStream` subclasses). The remaining items are
proportionate to an M0/draft-scope feature — a stale-branch CI signal, one open
reviewer disagreement that needs a definitive close-out rather than staying
unresolved, and a handful of non-blocking hardening/coverage suggestions. I
don't see a case for a fundamentally different architecture here; the
checkpoint-format-scoping and class-allowlist patterns are the right primitives
for this problem
. This posts as a comment rather than an approval/changes-requested review: it
is my own PR, and GitHub does not permit self-approval or self-request-changes
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]