eskabetxe commented on code in PR #4521:
URL: https://github.com/apache/flink-cdc/pull/4521#discussion_r4127390713


##########
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:
   Good catch, you're right — the 2.7 search guard only works because upstream 
had already removed the hasCompletelyProcessedPosition() guard around heartbeat 
dispatch in 2.4 (DBZ-6635, "Send heartbeats also before processing first 
event"). Since this fork is based on Debezium 1.9.8, that guard was still in 
place, so skipping the search would indeed just have moved the stall into 
processMessages().
   
   Pushed a follow-up that completes the backport:
   - processMessages() now dispatches heartbeats unconditionally when no 
message is received (DBZ-6635), so heartbeat.action.query can generate WAL on a 
fresh start;
   - searchWalPosition() also dispatches heartbeats while waiting, mirroring 
2.7, so the search terminates on an idle publication when resuming.
   
   Added a unit test that runs execute() with a fresh-start offset and a stream 
that never yields a message, asserting heartbeat dispatch happens (and the 
search is skipped); it fails without the fix. PTAL, thanks!



-- 
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]

Reply via email to