[
https://issues.apache.org/jira/browse/FLINK-40560?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40560:
-----------------------------------
Labels: pull-request-available (was: )
> [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
> Labels: pull-request-available
>
> 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}
> Those two offsets are not a completion signal. For `snapshot-only`,
> `HybridSplitAssigner.createStreamSplit()` initialises them from 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 it starts out as `min(high watermarks) >= max(high watermarks)`, which
> only holds when every
> snapshot split shares one high watermark — a single split, or a completely
> idle cluster.
> Afterwards the split's starting offset advances only through
> `IncrementalSourceRecordEmitter.updateStreamSplitState`, and only for
> data-change records and
> heartbeat events:
> {code:java}
> private void updateStreamSplitState(SourceSplitState splitState, SourceRecord
> element) {
> if (splitState.isStreamSplitState()) {
> Offset position = getOffsetPosition(element);
> splitState.asStreamSplitState().setStartingOffset(position);
> }
> }
> {code}
> So the guard really asks "have emitted records carried the split watermark
> past the stopping
> offset", which is a different question from "has the bounded stream split
> finished". A
> `snapshot-only` split reaches its stopping offset and completes without
> necessarily emitting any
> record that would move that watermark — a captured publication with no
> traffic emits none at all.
> `onSplitFinished` is invoked precisely because the split completed, so the
> offset comparison adds
> no safety; it just suppresses the cleanup in the common case.
> h3. Expected
> In `snapshot-only` mode the slot should be dropped once the bounded stream
> split has finished.
> h3. Suggested fix
> Key the cleanup off "the bounded stream split finished" instead of the offset
> comparison.
> `IncrementalSourceReader.onSplitFinished` already documents that a stream
> split finishes for
> exactly two reasons:
> {code:java}
> // Two possibilities that finish a stream split:
> //
> // 1. Stream reader is suspended by enumerator because new tables have been
> // finished its snapshot reading.
> // Under this case
> incrementalSourceReaderContext.isStreamSplitReaderSuspended()
> // is true and need to request the latest finished splits number.
> //
> // 2. Stream reader reaches the ending offset of the split. We need to do
> // nothing under this case.
> {code}
> So `isSnapshotOnly() &&
> !incrementalSourceReaderContext.isStreamSplitReaderSuspended()` expresses
> "the bounded split ran to its ending offset" directly, without depending on
> which records happened
> to be emitted. It also drops the offset dereference, removing a latent NPE
> when the stopping
> offset is null.
> Two notes for whoever picks this up:
> * `PostgresDialect.removeSlot` swallows exceptions and returns `false`, so
> today's
> `LOG.info("Remove slot '{}' result is {}.")` reports a failed cleanup at
> INFO. Worth raising to
> WARN when the result is `false`.
> * Dropping the slot from `close()` instead would be wrong: a job stopped for
> a later restore still
> needs its `restart_lsn`, so that trades a leaked slot for data loss.
> 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)