[
https://issues.apache.org/jira/browse/FLINK-40538?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
João Boto updated FLINK-40538:
------------------------------
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).
h3. What happened
A Postgres CDC source pointed at a database that receives no writes to its
publication
never leaves Debezium's WAL position search. It never begins streaming, and:
- no change events are ever produced (including once writes eventually start,
until the
first of them arrives);
- `heartbeat.action.query` never executes, so `heartbeat.interval.ms` has no
effect;
- the replication slot's `confirmed_flush_lsn` never advances, so PostgreSQL
retains
every WAL segment from the slot's position for as long as the job runs.
The last point is the damaging one: the job reports RUNNING and healthy while
WAL grows
without bound on the source database.
h3. How to reproduce
1. Create a Postgres CDC source over a schema whose tables receive no writes.
2. Start the job. Other databases in the cluster may be active; the captured
publication must be quiet.
3. Observe the job is RUNNING, all checkpoints complete, and `numRecordsOut` of
the
source is 0.
4. Observe `confirmed_flush_lsn` in `pg_replication_slots` stays fixed while
`pg_current_wal_lsn()` advances.
5. Write any row to a captured table — the connector immediately starts
streaming and
the slot begins advancing.
h3. Logs
At DEBUG for `io.debezium`, TRACE for
`io.debezium.connector.postgresql.connection`:
INFO WalPositionLocator - Looking for WAL restart position for last commit LSN
'null'
and last change LSN 'LSN\{0/518EB00}'
INFO PostgresStreamingChangeEventSource - Searching for WAL resume position
TRACE PostgresReplicationConnection - Streaming requested from LSN
LSN\{0/518EB00}, received LSN LSN\{0/5192A60}
TRACE PostgresReplicationConnection - Streaming requested from LSN
LSN\{0/518EB00}, received LSN LSN\{0/5192AF0}
... repeats indefinitely, requested LSN never changes ...
`"WAL resume position '{}' discovered"` is never logged.
h3. Cause
`PostgresSourceFetchTaskContext.loadStartingOffsetState` always returns a
non-null
`PostgresOffsetContext`, built from the stream split's starting offset. In
`PostgresStreamingChangeEventSource#execute` that makes
```java
boolean hasStartLsnStoredInContext = offsetContext != null;
always true, so the WAL position search branch is taken unconditionally. Plain
Debezium takes the other branch on a fresh start and logs
"WAL position will not be searched".
searchWalPosition then loops until a message is decoded, and — unlike the main
streaming loop — dispatches no heartbeat while waiting:
// searchWalPosition
if (receivedMessage) \{ noMessageIterations = 0; }
else {
noMessageIterations++; // no dispatchHeartbeatEvent
...
}
// main streaming loop
else {
if (offsetContext.hasCompletelyProcessedPosition()) {
dispatcher.dispatchHeartbeatEvent(partition, offsetContext);
}
...
}
That is a deadlock: the search waits for publication traffic, and the only
mechanism
that would generate traffic on an idle database — heartbeat.action.query,
driven by
the heartbeat dispatch — runs only from the main loop, which the search
precedes.
Possible fixes
1. Do not force the search when the stored offset has nothing processed yet, so
Debezium takes the new WalPositionLocator() path and starts from the slot's
confirmed position.
2. Write one WAL record when the stream task starts, so the search terminates
deterministically. We are running this as a local patch: a single
pg_logical_emit_message(false, '<internal prefix>', ...) in
PostgresStreamFetchTask#execute before delegating to
StreamSplitReadTask#execute. The prefix matches no configured logical
message, so
the record is filtered before the sink.
3. Upstream Debezium: dispatch heartbeats in searchWalPosition as the main loop
does.
This would fix it without a change here, but only for versions carrying that
fix.
h3. Are you willing to submit a PR?
Yes.
was:
### 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).
### What happened
A Postgres CDC source pointed at a database that receives no writes to its
publication
never leaves Debezium's WAL position search. It never begins streaming, and:
- no change events are ever produced (including once writes eventually start,
until the
first of them arrives);
- `heartbeat.action.query` never executes, so `heartbeat.interval.ms` has no
effect;
- the replication slot's `confirmed_flush_lsn` never advances, so PostgreSQL
retains
every WAL segment from the slot's position for as long as the job runs.
The last point is the damaging one: the job reports RUNNING and healthy while
WAL grows
without bound on the source database.
### How to reproduce
1. Create a Postgres CDC source over a schema whose tables receive no writes.
2. Start the job. Other databases in the cluster may be active; the captured
publication must be quiet.
3. Observe the job is RUNNING, all checkpoints complete, and `numRecordsOut` of
the
source is 0.
4. Observe `confirmed_flush_lsn` in `pg_replication_slots` stays fixed while
`pg_current_wal_lsn()` advances.
5. Write any row to a captured table — the connector immediately starts
streaming and
the slot begins advancing.
### Logs
At DEBUG for `io.debezium`, TRACE for
`io.debezium.connector.postgresql.connection`:
INFO WalPositionLocator - Looking for WAL restart position for last commit LSN
'null'
and last change LSN 'LSN\{0/518EB00}'
INFO PostgresStreamingChangeEventSource - Searching for WAL resume position
TRACE PostgresReplicationConnection - Streaming requested from LSN
LSN\{0/518EB00}, received LSN LSN\{0/5192A60}
TRACE PostgresReplicationConnection - Streaming requested from LSN
LSN\{0/518EB00}, received LSN LSN\{0/5192AF0}
... repeats indefinitely, requested LSN never changes ...
`"WAL resume position '{}' discovered"` is never logged.
### Cause
`PostgresSourceFetchTaskContext.loadStartingOffsetState` always returns a
non-null
`PostgresOffsetContext`, built from the stream split's starting offset. In
`PostgresStreamingChangeEventSource#execute` that makes
```java
boolean hasStartLsnStoredInContext = offsetContext != null;
always true, so the WAL position search branch is taken unconditionally. Plain
Debezium takes the other branch on a fresh start and logs
"WAL position will not be searched".
searchWalPosition then loops until a message is decoded, and — unlike the main
streaming loop — dispatches no heartbeat while waiting:
// searchWalPosition
if (receivedMessage) \{ noMessageIterations = 0; }
else {
noMessageIterations++; // no dispatchHeartbeatEvent
...
}
// main streaming loop
else {
if (offsetContext.hasCompletelyProcessedPosition()) {
dispatcher.dispatchHeartbeatEvent(partition, offsetContext);
}
...
}
That is a deadlock: the search waits for publication traffic, and the only
mechanism
that would generate traffic on an idle database — heartbeat.action.query,
driven by
the heartbeat dispatch — runs only from the main loop, which the search
precedes.
Possible fixes
1. Do not force the search when the stored offset has nothing processed yet, so
Debezium takes the new WalPositionLocator() path and starts from the slot's
confirmed position.
2. Write one WAL record when the stream task starts, so the search terminates
deterministically. We are running this as a local patch: a single
pg_logical_emit_message(false, '<internal prefix>', ...) in
PostgresStreamFetchTask#execute before delegating to
StreamSplitReadTask#execute. The prefix matches no configured logical
message, so
the record is filtered before the sink.
3. Upstream Debezium: dispatch heartbeats in searchWalPosition as the main loop
does.
This would fix it without a change here, but only for versions carrying that
fix.
### Are you willing to submit a PR?
Yes.
> [Postgres] Connector never starts streaming on a database with no publication
> traffic
> -------------------------------------------------------------------------------------
>
> Key: FLINK-40538
> URL: https://issues.apache.org/jira/browse/FLINK-40538
> Project: Flink
> Issue Type: Bug
> Components: Flink CDC
> Affects Versions: cdc-3.6.0
> Reporter: João Boto
> Priority: Major
> Fix For: cdc-3.6.0
>
>
> 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).
> h3. What happened
> A Postgres CDC source pointed at a database that receives no writes to its
> publication
> never leaves Debezium's WAL position search. It never begins streaming, and:
> - no change events are ever produced (including once writes eventually start,
> until the
> first of them arrives);
> - `heartbeat.action.query` never executes, so `heartbeat.interval.ms` has no
> effect;
> - the replication slot's `confirmed_flush_lsn` never advances, so PostgreSQL
> retains
> every WAL segment from the slot's position for as long as the job runs.
> The last point is the damaging one: the job reports RUNNING and healthy while
> WAL grows
> without bound on the source database.
> h3. How to reproduce
> 1. Create a Postgres CDC source over a schema whose tables receive no writes.
> 2. Start the job. Other databases in the cluster may be active; the captured
> publication must be quiet.
> 3. Observe the job is RUNNING, all checkpoints complete, and `numRecordsOut`
> of the
> source is 0.
> 4. Observe `confirmed_flush_lsn` in `pg_replication_slots` stays fixed while
> `pg_current_wal_lsn()` advances.
> 5. Write any row to a captured table — the connector immediately starts
> streaming and
> the slot begins advancing.
> h3. Logs
> At DEBUG for `io.debezium`, TRACE for
> `io.debezium.connector.postgresql.connection`:
> INFO WalPositionLocator - Looking for WAL restart position for last commit
> LSN 'null'
> and last change LSN 'LSN\{0/518EB00}'
> INFO PostgresStreamingChangeEventSource - Searching for WAL resume position
> TRACE PostgresReplicationConnection - Streaming requested from LSN
> LSN\{0/518EB00}, received LSN LSN\{0/5192A60}
> TRACE PostgresReplicationConnection - Streaming requested from LSN
> LSN\{0/518EB00}, received LSN LSN\{0/5192AF0}
> ... repeats indefinitely, requested LSN never changes ...
> `"WAL resume position '{}' discovered"` is never logged.
> h3. Cause
> `PostgresSourceFetchTaskContext.loadStartingOffsetState` always returns a
> non-null
> `PostgresOffsetContext`, built from the stream split's starting offset. In
> `PostgresStreamingChangeEventSource#execute` that makes
> ```java
> boolean hasStartLsnStoredInContext = offsetContext != null;
> always true, so the WAL position search branch is taken unconditionally. Plain
> Debezium takes the other branch on a fresh start and logs
> "WAL position will not be searched".
> searchWalPosition then loops until a message is decoded, and — unlike the main
> streaming loop — dispatches no heartbeat while waiting:
> // searchWalPosition
> if (receivedMessage) \{ noMessageIterations = 0; }
> else {
> noMessageIterations++; // no dispatchHeartbeatEvent
> ...
> }
> // main streaming loop
> else {
> if (offsetContext.hasCompletelyProcessedPosition()) {
> dispatcher.dispatchHeartbeatEvent(partition, offsetContext);
> }
> ...
> }
> That is a deadlock: the search waits for publication traffic, and the only
> mechanism
> that would generate traffic on an idle database — heartbeat.action.query,
> driven by
> the heartbeat dispatch — runs only from the main loop, which the search
> precedes.
> Possible fixes
> 1. Do not force the search when the stored offset has nothing processed yet,
> so
> Debezium takes the new WalPositionLocator() path and starts from the slot's
> confirmed position.
> 2. Write one WAL record when the stream task starts, so the search terminates
> deterministically. We are running this as a local patch: a single
> pg_logical_emit_message(false, '<internal prefix>', ...) in
> PostgresStreamFetchTask#execute before delegating to
> StreamSplitReadTask#execute. The prefix matches no configured logical
> message, so
> the record is filtered before the sink.
> 3. Upstream Debezium: dispatch heartbeats in searchWalPosition as the main
> loop does.
> This would fix it without a change here, but only for versions carrying
> that fix.
> h3. Are you willing to submit a PR?
> Yes.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)