DanielLeens commented on PR #11614:
URL: https://github.com/apache/seatunnel/pull/11614#issuecomment-5379375339
Thanks for keeping this moving. I want to open with a correction: on
2026-08-01 I reviewed this PR (COMMENTED) and again on 2026-08-20 I re-reviewed
the current head and APPROVED it as "Fully compatible" with "no source-level
blocker." Doing a completely fresh, independent pass today, I found a real
compatibility gap that both of my earlier rounds missed. I'm sorry for
approving past it — here's the honest re-analysis.
# What Problem Does This PR Solve?
Before this PR, a sink writer could inherit the deprecated, default no-op
`SinkWriter.applySchemaChange(SchemaChangeEvent)` method. If it then received a
CDC `SchemaChangeEvent` (e.g. a source-side `ALTER TABLE ADD COLUMN`), the
event was silently swallowed: no exception, no log warning at the call site,
nothing — the job just kept running while its sink schema quietly drifted out
of sync with the incoming data.
This PR closes that gap by routing every `applySchemaChange` call through a
single new static helper,
`SupportSchemaEvolutionSinkWriter.applySchemaChangeToWriter(...)`, which allows
the call only when the writer either (a) implements
`SupportSchemaEvolutionSinkWriter`, or (b) genuinely overrides the deprecated
method (detected via reflection). Any other writer now gets an explicit
`UnsupportedOperationException` instead of a silent no-op. Iceberg's writer is
also changed to throw explicitly when
`iceberg.table.schema-evolution-enabled=false` instead of silently ignoring the
event, and Lance no longer falsely advertises schema-evolution support via an
empty override.
One-sentence summary: this PR turns a previously silent, hard-to-detect
schema-drift failure mode into an explicit, fail-fast exception across Zeta,
Flink, and multi-table sink dispatch.
# 1. Code Change Review
## 1.1 Core Logic Analysis
I independently re-verified the diff against the current PR head
(`61be2cf60d75141a06817c2128eda3165e72e02c`) via `git diff` against the
merge-base with `dev` — 12 files, matching what both of my prior rounds
reported file-for-file.
Core logic, read directly from
`seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/SupportSchemaEvolutionSinkWriter.java`:
```java
static void applySchemaChangeToWriter(SinkWriter<?, ?, ?> writer,
SchemaChangeEvent event)
throws IOException {
if (writer instanceof SupportSchemaEvolutionSinkWriter) {
((SupportSchemaEvolutionSinkWriter) writer).applySchemaChange(event);
return;
}
if (overridesDeprecatedApplySchemaChange(writer)) {
writer.applySchemaChange(event);
return;
}
throw new UnsupportedOperationException(...);
}
```
`overridesDeprecatedApplySchemaChange` uses
`writer.getClass().getMethod("applySchemaChange",
SchemaChangeEvent.class).getDeclaringClass() != SinkWriter.class` to detect a
real legacy override. Since the interface method is implicitly public, any
override must also be public, so `getMethod` (public-only lookup) is sufficient
here — I don't see a case where a legitimate override would be missed. This
runs once per schema-change event (a rare, DDL-driven event), not per record,
so there's no hot-path cost.
**Simple before/after example**, traced myself against the real call chain
rather than assumed: say a CDC pipeline writes MySQL into Iceberg, and the user
never sets `iceberg.table.schema-evolution-enabled` (documented default is
`false` — confirmed in `IcebergSinkOptions.java:47` and
`docs/en/connectors/sink/Iceberg.md:76`). A DBA runs `ALTER TABLE ... ADD
COLUMN` on the source.
- **Before this PR**: `IcebergSinkWriter.applySchemaChange` just returns
early when the option is `false`. The job keeps running; the new column's data
is effectively dropped downstream. Silent, but alive.
- **After this PR**: the same event now throws
`UnsupportedOperationException("Iceberg sink received schema change event for
table %s, but iceberg.table.schema-evolution-enabled is false.")`.
I traced where that exception actually goes, since that determines
real-world blast radius:
```
Zeta direct (non-multi-table) sink path
CDC source emits SchemaChangeEvent
-> SinkFlowLifeCycle.received(record)
(SinkFlowLifeCycle.java:268)
-> record.getData() instanceof SchemaChangeEvent
-> processSchemaChangeEvent(event)
(SinkFlowLifeCycle.java:825)
->
SupportSchemaEvolutionSinkWriter.applySchemaChangeToWriter(writer, event)
-> unsupported writer -> throw UnsupportedOperationException
-> received()'s own catch block:
catch (RuntimeException e) { throw e; }
(SinkFlowLifeCycle.java:291-292)
-> propagates out of received() uncaught -> sink task fails -> job fails
```
I read `received()` directly (lines 268-293) to confirm this: it catches
`RuntimeException` and rethrows it as-is, and wraps any other `Exception` in a
new `RuntimeException`. There is no "log and continue" branch for a
schema-change failure at this layer. So for a single-table sink using the
default config, this is a genuine job-crashing change, not a cosmetic one.
For **multi-table** sinks, I also checked
`MultiTableSinkWriter.applySchemaChange` — the diff only swaps the old inline
instanceof-check for the same shared helper inside the pre-existing
`executeWithTableRetry(..., MultiTableFailurePhase.CHECKPOINT, ...)` wrapper,
so the pre-existing `MultiTableFailurePolicy` (FAIL_FAST vs.
CONTINUE_OTHER_TABLES) still governs blast radius here — that dispatch logic
itself isn't touched by this PR. Worth knowing: I checked
`MultiTableSinkWriter.java:126`, and one of the writer's constructor overloads
defaults `failurePolicy` to `MultiTableFailurePolicy.FAIL_FAST`, so unless a
job explicitly configured `CONTINUE_OTHER_TABLES`, the multi-table path also
fails the whole sink, not just the affected table.
For **Flink**, I read the full current
`FlinkSinkWriter.handleSchemaChangeEvent` (both the `-common` and `-20` copies
are structurally identical). This is where I found something genuinely good
that neither of my prior rounds called out: the **old** code returned early on
an unsupported writer without ever calling `sendSchemaChangeAck(...)`:
```java
// before
if (!(sinkWriter instanceof SupportSchemaEvolutionSinkWriter)) {
log.warn(...);
return; // no ack sent to LocalSchemaCoordinator at all
}
```
The **new** code always goes through a `try { applySchemaChangeToWriter(...)
} catch (Exception e) { log.error(...) } finally { sendSchemaChangeAck(...,
success) }` block, so the coordinator is now notified of failure even for
unsupported writers, instead of potentially waiting forever for an ack that the
old code path would never send. That's a real, independently-verified
correctness improvement on top of the headline fail-fast change, and it's not
mentioned in either of my earlier reviews.
**Key findings:**
- The single dispatch choke point is correctly wired into all five call
sites (`SinkFlowLifeCycle`, `MultiTableSinkWriter`, `ErrorHandlingSinkWriter`,
and both `FlinkSinkWriter` copies) — I read each one directly, not just the
diff hunks.
- The reflection-based legacy-override detection is correct and bounded to a
rare event type, not the record hot path.
- The Flink wrapper's ack-on-failure fix (via the `finally` block) closes a
real coordinator-hang risk that existed before this PR, independent of the
headline fail-fast behavior.
- The actual production-facing risk of this PR is not in its logic — it's
that the resulting behavior change is a real backward-incompatibility for a
specific, common historical configuration, and it isn't documented as one. See
1.2.
## 1.2 Compatibility Impact
**Partially incompatible** — this is a correction to my prior two rounds,
both of which said "Fully compatible."
No API, SPI method, config option, or serialization format is removed or
renamed, and `iceberg.table.schema-evolution-enabled`'s default value is
unchanged (`false`). So in the narrow "did you rename/remove a contract" sense,
my earlier assessment wasn't wrong.
But this project's own `CLAUDE.md` treats backward compatibility as covering
"historical behavior," not just the config/API surface, and explicitly
requires: *"Any incompatible change MUST be explicitly documented... in
`docs/en/introduction/concepts/incompatible-changes.md`... include migration
guidance... be clearly explained in the PR description."* I checked
`docs/en/introduction/concepts/incompatible-changes.md` at this PR's head: it
already has a dedicated (currently empty) `### Engine Behavior Changes`
section, and this PR's change is exactly the kind of entry that section exists
for — but no entry was added, and the PR description doesn't flag this as a
breaking change either (it frames it purely as closing a
documentation-vs-runtime gap).
The concrete production impact: any currently-running job that (a) uses a
sink writer that doesn't implement `SupportSchemaEvolutionSinkWriter` and
doesn't override the deprecated method (this includes Iceberg with the default
`schema-evolution-enabled=false`, and any third-party/custom sink not on the
small supported list), and (b) has a source that ever emits a
`SchemaChangeEvent` — will go from "keeps running with a silently stale schema"
to "job fails" after upgrading to this version, with **zero config or code
change on the user's side**, and **no way to opt back into the old lenient
behavior even temporarily** (I looked for an escape hatch — there isn't one).
That's the right long-term direction (silent data drift is worse than a loud
failure), but it's a real "surprise on upgrade" that this project's own policy
says must be documented with migration guidance, and currently isn't.
## 1.3 Performance / Side-Effect Analysis
No new thread pools, executors, connections, or locks. The added reflection
call in `overridesDeprecatedApplySchemaChange` runs once per schema-change
event (DDL-driven, rare), not per record — no measurable hot-path cost. No new
resource-release concerns; `close()`/writer lifecycle paths are untouched by
this diff.
## 1.4 Error Handling and Logging
No sensitive data in the new exception messages — they include the writer's
class name and the table path, which is exactly the context needed to diagnose
a fail-fast rejection, and nothing credential-like. See Issue 1 below for the
one substantive gap I found (documentation of the behavior change, not the
error-handling code itself, which is sound).
# 2. Code Quality Assessment
## 2.1 Coding Standards
AOSP formatting, ASF license headers present on the new test file, no
wildcard imports. The two new static methods on
`SupportSchemaEvolutionSinkWriter` both carry proper multi-line Javadoc
explaining intent and the legacy-compat rationale — no missing-comment gap here.
## 2.2 Test Coverage and Test Stability
I read the full content of all four new/changed test files in the diff
(`SupportSchemaEvolutionSinkWriterTest` — 152 lines, new;
`IcebergSinkWriterTest` — +29 lines; `ErrorHandlingSinkWriterTest` — +23 lines;
`FlinkSinkWriterTest` — +36 lines). All are pure JUnit 5 + Mockito unit tests:
no `Thread.sleep`, no Docker/network/fixed ports, no shared static mutable
state between tests, no order dependence, deterministic assertions (exact
exception type + message-contains checks). **Stability rating: Stable.**
Coverage is good for the headline behavior (supported-writer pass-through,
legacy-override pass-through, unsupported-writer fail-fast, Iceberg-disabled
fail-fast, Flink unsupported-writer failure-ack-then-exception) — I confirmed
each of these by reading the actual assertions, not just the test names. One
non-blocking gap, carried over from my own 2026-08-01 review and still true
today: `MultiTableSinkWriter`'s unsupported-sub-writer fail-fast path (the diff
at `MultiTableSinkWriter.java`) has no dedicated unit test of its own; it's
only exercised indirectly through the shared helper's own tests. I wouldn't
block on this — the shared helper is the actual logic under test, and the
multi-table wiring is a thin pass-through — but it would be a good addition.
## 2.3 Documentation Updates
This is where Issue 1 actually lives (see 1.2):
`docs/en/introduction/concepts/incompatible-changes.md` needs a new entry under
its existing `### Engine Behavior Changes` heading, in the same
`**[BREAKING]**` / Impact / Migration Guide format the file already uses for
other entries (I read several existing entries in that file to confirm the
expected format). `docs/en/connectors/sink/Iceberg.md`'s description of
`iceberg.table.schema-evolution-enabled` (currently just "Setting to true
enables Iceberg tables to support schema evolution...") should also gain a
sentence on what happens when it's `false` and a schema change event arrives
now (job failure, not silent ignore). The `docs/zh` counterparts need the same
updates per this repo's own contribution guide.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
**Precise fix** for the stated problem. Centralizing the decision in one
static choke point rather than duplicating the implements-check at each of the
five call sites is the right shape — there's exactly one place to audit, and I
confirmed all five call sites actually route through it rather than
re-implementing the check locally.
## 3.2 Maintainability
Improved. New sink writers only need to decide whether to implement
`SupportSchemaEvolutionSinkWriter`; the fail-fast behavior is automatic and
centrally enforced rather than something every engine integration has to
remember to re-check.
## 3.3 Extensibility
Good — the interface-based capability check composes cleanly with
`MultiTableSinkWriter`'s existing `FAIL_FAST`/`CONTINUE_OTHER_TABLES` per-table
policy without new plumbing, and with the Flink coordinator's existing ack
protocol.
## 3.4 Historical-Version Compatibility
No persisted checkpoint/savepoint format, state schema, or serialization
contract is touched. The compatibility concern here (see 1.2) is about *runtime
behavior* for historically-configured jobs on upgrade, not about state-restore
or serialization compatibility — those remain unaffected.
# 4. Issue Summary
| # | Issue | Location | Severity |
|---|---|---|---|
| 1 | Fail-fast behavior for sinks that previously silently no-op'd a schema
change (including any sink with `iceberg.table.schema-evolution-enabled=false`,
the documented default) is a real behavior change for existing production jobs
on upgrade, with no opt-out, but is not documented as an incompatible/breaking
change in `docs/en/introduction/concepts/incompatible-changes.md` (which has an
empty `Engine Behavior Changes` section ready for exactly this), nor called out
as a migration item in the PR description or in
`docs/en/connectors/sink/Iceberg.md`'s option description |
`seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/SupportSchemaEvolutionSinkWriter.java`;
`seatunnel-connectors-v2/connector-iceberg/.../IcebergSinkWriter.java:118-131`;
`docs/en/introduction/concepts/incompatible-changes.md` (missing entry) |
Medium |
| 2 | `MultiTableSinkWriter`'s unsupported-sub-writer fail-fast path has no
dedicated unit test |
`seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java`
| Low |
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
**1. Blockers — must be fixed:**
- Issue 1 (Medium): add the missing incompatible-change documentation entry
(and the accompanying Iceberg doc-option sentence) before this leaves draft.
This is a project policy requirement (`CLAUDE.md`'s backward-compatibility
rules), not a code defect — the underlying behavior change itself is the right
direction, it just needs to be written down with migration guidance so
operators aren't surprised by a job failing after a routine upgrade.
- CI: live-checked just now against the fork (`danielnadean/seatunnel`) run
for head `61be2cf60`: 76 success, 10 skipped, and 3 non-success —
`paimon-connector-it (8, ubuntu-latest)` failed, `all-connectors-it-7 (8,
ubuntu-latest)` failed, `kudu-connector-it (11, ubuntu-latest)` cancelled. None
of those modules line up with what this diff touches (`seatunnel-api`,
`seatunnel-engine-server`, `connector-iceberg`, `connector-lance`,
`seatunnel-translation-flink`), which is consistent with the unrelated-IT-flake
pattern seen in earlier rounds of this PR (rocketmq, then couchbase). I want to
flag honestly, though: I was not able to pull the actual failure logs this
round (the Actions log-streaming API returned a local stream error both times I
tried), so this diagnosis is based on module/job-name correlation rather than a
confirmed root-cause read — weaker evidence than my earlier rounds. Please
paste the failing job link(s) here, or just resync and rerun (see next point)
and I'll
take another look.
- Branch sync: live-checked via `gh api
repos/apache/seatunnel/compare/dev...61be2cf60...` just now — still
`ahead_by=1`, `behind_by=16`, `diverged`. That's worse than what my 2026-08-20
review reported ("zero commits behind") — `dev` has moved on since then and
this branch hasn't been resynced. Please sync with the latest `dev` and rerun
CI; that will also refresh the two IT lanes above in case they've already been
fixed upstream.
**2. Recommended fixes — non-blocking:**
- Issue 2 (Low): add a targeted unit test for `MultiTableSinkWriter`'s
unsupported sub-writer fail-fast path.
**Overall assessment.** The core approach is sound and I independently
re-verified it end to end rather than trusting my own prior conclusion: one
shared dispatch choke point, correctly wired into all five call sites, with a
genuine bonus fix in the Flink wrapper (failure acks are now always sent,
closing what looks like a real coordinator-hang gap in the old code). The tests
are stable and cover the headline behavior well. The one real gap is that this
PR changes historical runtime behavior for a common configuration (any sink
without explicit schema-evolution support, including Iceberg's own documented
default) without the incompatible-change documentation this project's own rules
require — that's a quick, doc-only fix, not a sign the implementation is wrong,
and I'd be glad to take another look as soon as it's added and CI is green
against a resynced `dev`.
--
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]