Rui Fan created FLINK-40517:
-------------------------------
Summary: Checkpointing-during-recovery drops cross-channel
watermark / watermark-status aggregation on fan-in rescale
Key: FLINK-40517
URL: https://issues.apache.org/jira/browse/FLINK-40517
Project: Flink
Issue Type: Sub-task
Components: Runtime / Checkpointing
Reporter: Rui Fan
Assignee: Rui Fan
On the normal recovery path, when a rescale fans several old input channels
into one new channel,
[`DemultiplexingRecordDeserializer.getNextRecord`](https://github.com/apache/flink/blob/76e774e04ff84ea305714153c5f4790fff98437c/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/io/recovery/DemultiplexingRecordDeserializer.java#L108-L131)
aggregates the recovered non-record elements across the merged old channels:
it emits the **min** watermark (suppressing it until every merged channel has
one) and `ACTIVE` if **any** channel is active. This reproduces the pre-rescale
`min(wmA, wmB, ...)` semantics so the operator watermark is never advanced past
what any old channel justifies.
With *checkpointing during recovery* (CDR,
`execution.checkpointing.during-recovery.enabled`) enabled, this aggregation
moves to the unspilling thread, but
[`GateFilterHandler.filterAndRewrite`](https://github.com/apache/flink/blob/76e774e04ff84ea305714153c5f4790fff98437c/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/ChannelStateFilteringHandler.java#L294-L302)
only filters records and passes watermarks and statuses through **verbatim**.
So on a fan-in rescale a high watermark from one merged channel advances the
operator watermark prematurely and legitimately in-flight lower-timestamp
records get dropped as late (silent data loss), and an `IDLE` from one merged
channel wrongly marks the whole channel idle (broken alignment; no data loss).
**Trigger**: CDR enabled + unaligned-checkpoint restore with in-flight state +
a fan-in rescale (scale-down / ambiguous keyed exchange); the data-loss symptom
additionally needs event-time watermarks and a late-dropping operator. Without
CDR the same rescale is correct.
**Fix**: do the same min-watermark / any-active aggregation in the unspilling
thread, where the per-old-channel `VirtualChannel`s already hold
`lastWatermark` / `watermarkStatus`, instead of passing them through verbatim
(`NO_RESCALE` needs no change). Since segments are written per new channel in
arrival order, a watermark must be held until all old channels folding into
that new channel are drained.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)