Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,15 @@ public void processElement2(StreamRecord<EqualityConvertPlan> record) {

@Override
public void processWatermark(Watermark mark) throws Exception {
if (planResult == null || mark.getTimestamp() < planResult.doneTimestamp()) {
if (planResult == null) {
// No conversion cycle is active: forward the watermark. It may originate from another
// task's trigger sharing the maintenance lock, and holding it back would stall lock
// release for that task.
super.processWatermark(mark);
return;
}

if (mark.getTimestamp() < planResult.doneTimestamp()) {
// Hold back watermarks until the cycle commits so the LockRemover keeps the maintenance lock
// for the whole cycle. Forwarding the planner's mid-cycle phase watermarks will release the
// lock early and could let the next trigger run a concurrent cycle on the same staging
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,24 @@ void holdsBackWatermarkUntilCommit() throws Exception {
}
}

@Test
void forwardsWatermarkWithoutActiveCycle() throws Exception {
Table table = createTable(3, FileFormat.PARQUET);
insert(table, 1, "a");

try (TwoInputStreamOperatorTestHarness<DVWriteResult, EqualityConvertPlan, Trigger> harness =
createHarness()) {
harness.open();

// No plan was received, so no conversion cycle is active. A watermark here originates
// from another task's trigger sharing the maintenance lock; dropping it would stall
// the unioned watermark and the lock would never be released.
long markTs = System.currentTimeMillis();
harness.processBothWatermarks(new Watermark(markTs));
assertThat(watermarks(harness)).containsExactly(new Watermark(markTs));
}
}

@Test
void skipsCommitForEmptyCycle() throws Exception {
Table table = createTable(3, FileFormat.PARQUET);
Expand Down
Loading