[
https://issues.apache.org/jira/browse/FLINK-40560?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18111781#comment-18111781
]
tchivs commented on FLINK-40560:
--------------------------------
Corrected the Cause section. The earlier wording said the guard can only hold
when all snapshot splits share one high watermark, which is only half of it:
the split's starting offset is also advanced afterwards by
IncrementalSourceRecordEmitter.updateStreamSplitState, but only for data-change
records and heartbeat events. So the guard is asking whether emitted records
have carried the split watermark past the stopping offset, which is a different
question from whether the bounded split finished — and a snapshot-only split
over a quiet publication completes without emitting anything that would move it.
> [Postgres] snapshot-only replication slot is never dropped unless the
> snapshot is instantaneous
> -----------------------------------------------------------------------------------------------
>
> Key: FLINK-40560
> URL: https://issues.apache.org/jira/browse/FLINK-40560
> Project: Flink
> Issue Type: Bug
> Components: Flink CDC
> Affects Versions: cdc-3.6.0
> Reporter: tchivs
> Priority: Major
>
> h3. Flink version / Flink CDC version
> Flink 2.2.1, Flink CDC 3.6.0 (`flink-connector-postgres-cdc`), Debezium
> 1.9.8.Final,
> PostgreSQL 18 (not version-specific), `pgoutput`.
> h3. What happened
> FLINK-38277 added slot cleanup for `snapshot-only` sources: when the stream
> split finishes,
> `PostgresSourceReader.onSplitFinished` drops the replication slot. In
> practice the drop never
> happens unless the whole snapshot is effectively instantaneous, so every
> bounded snapshot run
> leaves its replication slot behind.
> The job reaches FINISHED normally, no `Remove slot '{}' result is {}.` line
> is ever logged, and
> the slot stays in `pg_replication_slots` with `active = false`, pinning WAL
> from its
> `restart_lsn` until an operator drops it manually.
> h3. How to reproduce
> # Configure a Postgres incremental source with `scan.startup.mode=snapshot`.
> # Capture enough tables/partitions that the snapshot takes more than a
> moment. Our case is six
> tables, four of them monthly-partitioned, expanding to ~380 snapshot splits
> over ~17 minutes.
> Any workload with concurrent writes elsewhere in the cluster is enough.
> # Let the job run to FINISHED.
> # Inspect `pg_replication_slots`: the slot is still there.
> h3. Logs
> The stream split completes and the reader reaches an offset past the stopping
> offset, yet the
> drop is skipped:
> {noformat}
> StreamSplitReadTask finished for StreamSplit{splitId='stream-split',
> offset=Offset{lsn=LSN{2A/385B31F8}, ...}, <- startingOffset = min(high
> watermarks)
> endOffset=Offset{lsn=LSN{2A/3D86A000}, ...}, <- endingOffset = max(high
> watermarks)
> isSnapshotCompleted=true}
> at Offset{lsn=LSN{2A/3D8E96A8}, ...}
> {noformat}
> No `Remove slot` line follows. On an otherwise idle development database, a
> single ~17 minute run
> retained 85-135 MB of WAL, and the slot keeps accumulating until it is
> dropped by hand.
> h3. Cause
> `PostgresSourceReader.onSplitFinished` guards the drop with an offset
> comparison:
> {code:java}
> if (this.sourceConfig.getStartupOptions().isSnapshotOnly()
> &&
> streamSplit.getStartingOffset().isAtOrAfter(streamSplit.getEndingOffset())) {
> boolean removed = dialect.removeSlot(dialect.getSlotName());
> LOG.info("Remove slot '{}' result is {}.", dialect.getSlotName(),
> removed);
> }
> {code}
> But `HybridSplitAssigner.createStreamSplit()` fills those two offsets with
> the *lowest* and the
> *highest* high-watermark observed across the finished snapshot splits:
> {code:java}
> // minOffset / maxOffset are the lowest / highest high watermark among
> finished snapshot splits
> Offset stoppingOffset = offsetFactory.createNoStoppingOffset();
> if (sourceConfig.getStartupOptions().isSnapshotOnly()) {
> stoppingOffset = maxOffset;
> }
> return new StreamSplit(
> STREAM_SPLIT_ID,
> minOffset == null ? offsetFactory.createInitialOffset() : minOffset,
> stoppingOffset,
> ...);
> {code}
> So the guard reduces to `min(high watermarks) >= max(high watermarks)`, which
> can only hold when
> every snapshot split carries the exact same high watermark — that is, when no
> WAL at all is
> produced between the first and the last split finishing. With more than one
> split, or any
> concurrent write activity in the cluster, the condition is false by
> construction and the cleanup
> never runs. This is also why it passes in tests over small tables.
> h3. Expected
> In `snapshot-only` mode the slot should be dropped once the bounded stream
> split has finished.
> h3. Suggested fix
> Key the drop off "the bounded stream split finished" instead of
> `startingOffset.isAtOrAfter(endingOffset)`. `onSplitFinished` is only invoked
> for a split that
> already completed, so for `isSnapshotOnly()` the offset comparison adds no
> safety — it only
> suppresses the cleanup in every non-trivial case.
> h3. Workaround
> Registering a `JobStatusHook` that drops the slot on FINISHED / FAILED /
> CANCELED works. One
> caveat for anyone doing the same: a `Throwable` escaping a `JobStatusHook` is
> routed to Flink's
> `FatalExitExceptionHandler`, which terminates the JobManager process, so the
> hook has to contain
> every `Throwable` itself. Its JDBC class graph also needs to be warmed up
> before the terminal
> callback, because the user class loader is being torn down by then.
> h3. Related
> * FLINK-38277 — introduced the cleanup this issue reports as ineffective
> (fixed in cdc-3.5.0).
> * FLINK-40538 — separate bug in the same area; on a publication with no
> traffic the stream split
> cannot reach its stopping offset at all.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)