tchivs created FLINK-40560:
------------------------------
Summary: [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
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)