DanielLeens commented on issue #7937: URL: https://github.com/apache/seatunnel/issues/7937#issuecomment-5098190602
Dynamic Lookup implementation draft has been pushed to PR #11559. PR: https://github.com/apache/seatunnel/pull/11559 Head commit: `183bfbcc98aa3f5061ee8df7968d3f1c1a2e3d4a` Design-to-implementation summary: 1. Dedicated action and parser - Adds `dynamic_lookup` parser support. - Builds a dedicated `DynamicLookupAction` instead of modeling lookup as a transform. - Produces an output `CatalogTable` from the declared fact/dimension projection. 2. Fact source gate and restore ownership - Adds `FactSourceGateCapability`, `SourceGateState`, and `SourceGateCommand` API contracts. - Kafka source now advertises `FACT_SOURCE_GATE_V1` and stages splits while the fact gate is closed. - Closed-gate Kafka split ownership is snapshotted through native `snapshotState`, so restore keeps the source as the split owner. - Fact split activation happens only after the lookup side sends an explicit `OPEN` command. 3. Bootstrap and fact activation ordering - Dynamic lookup consumes dimension and fact ports from same-task-group blocking queues for this draft. - Dimension state is checkpointed first. - The first fact-gate `OPEN` is sent only from `notifyCheckpointComplete` for the checkpoint that contains the dynamic lookup dimension state. - If restoring from a durable lookup checkpoint, the fact gate opens during runtime open because the committed checkpoint already implies `FACT_POSITIONS_DURABLE`. 4. Checkpoint persistence and durable phase metadata - Adds `CheckpointIntent` to `CompletedCheckpoint`. - Dynamic lookup checkpoint contents are detected by state envelope magic. - Completed checkpoints carrying lookup state get purpose `DYNAMIC_LOOKUP_FACT_POSITION_ANCHOR`, target phase `FACT_POSITIONS_DURABLE`, and a SHA-256 digest over canonical action state bytes. - Completed checkpoint payloads are encoded with magic/version/payload length/payload/SHA-256 and strict legacy fallback. 5. Runtime lookup semantics - Maintains dimension state keyed by planner-resolved key fields. - Supports LEFT and INNER lookup projection. - Fact input is append-only in M1. - Dimension UPDATE_BEFORE/UPDATE_AFTER must be adjacent and same-key; different keys fail fast as primary-key updates. - Schema change events fail fast. - Checkpoint barriers align across fact and dimension ports before snapshot and downstream forwarding. 6. Channel and identity groundwork - Logical channel keys have canonical bytes and SHA-256 digest. - Channel attempt identity includes job execution epoch. - Channel envelopes include canonical digest and duplicate/conflict detection. - The current draft runtime uses same-task-group blocking queues; remote channel transport wiring remains a follow-up evidence gate before leaving Draft. 7. Resource admission - Parser enforces M1 mode/resource constraints: - disk-backed state - no TTL - one concurrent snapshot - max 4 GiB logical state per subtask - max 512 MiB resident state per subtask - It also checks the uncommitted snapshot count/byte formulas, cleanup throughput, local disk reservation, and remote staging quota before accepting the plan. Local checks performed: - `git diff --check` - `git diff --cached --check` - AOSP google-java-format apply/dry-run on the known changed Java files Not run locally by SeaTunnel local policy: - Maven compile - unit tests - integration tests - E2E tests - Docker/service startup - real SeaTunnel job execution The PR is intentionally still Draft. Current-head GitHub CI and maintainer review remain the compile/runtime evidence gate. -- 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]
