[ 
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)

Reply via email to