loserwang1024 commented on code in PR #4521:
URL: https://github.com/apache/flink-cdc/pull/4521#discussion_r4119355708
##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSource.java:
##########
@@ -187,7 +187,15 @@ public void execute(
this.lastCompletelyProcessedLsn =
replicationStream.get().startLsn();
- if (walPosition.searchingEnabled()) {
+ // Only search for the WAL resume position when the stored offset
has actually
+ // processed a position. On a fresh start (nothing processed yet)
the search loop
+ // would block forever waiting for a decoded message: on an idle
publication no WAL
+ // is produced, and the heartbeat action query that would generate
some only runs
+ // from the main streaming loop, which this search precedes.
Skipping the search here
+ // still starts streaming from the stored LSN, so no events are
missed. This mirrors
+ // the fix Debezium shipped in 2.7, which added the
hasCompletelyProcessedPosition()
+ // guard to searchingEnabled().
+ if (walPosition.searchingEnabled() &&
offsetContext.hasCompletelyProcessedPosition()) {
Review Comment:
Skipping the WAL-position search does not appear to resolve the
idle-publication issue, because `processMessages()` still guards heartbeat
dispatch with the same condition:
```java
if (offsetContext.hasCompletelyProcessedPosition()) {
dispatcher.dispatchHeartbeatEvent(partition, offsetContext);
}
```
On a fresh start, the offset contains a starting LSN but no `lsn_proc`, so
`hasCompletelyProcessedPosition()` returns `false`. Even after this change
allows execution to reach the main streaming loop, heartbeat dispatch is still
skipped.
This creates a circular dependency: processing a message is required to
enable heartbeats, but on an idle publication, `heartbeat.action.query` may be
the only mechanism that would generate that message. The query is invoked
through heartbeat dispatch; it does not run independently on a timer.
As a result, this change moves the wait from the WAL-position search into
the main streaming loop, while the heartbeat query remains blocked until an
external write produces a message that advances the processed position.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]