[
https://issues.apache.org/jira/browse/FLINK-40560?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
tchivs updated FLINK-40560:
---------------------------
Description:
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 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.
was:
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.
> [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}
> 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 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)