diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSource.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSource.java index e309a3b29a4..2ccf610a7b4 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSource.java +++ b/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()) { searchWalPosition(context, stream, walPosition); try { if (!isInPreSnapshotCatchUpStreaming(offsetContext)) { diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSourceTest.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSourceTest.java new file mode 100644 index 00000000000..993e199ec0d --- /dev/null +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSourceTest.java @@ -0,0 +1,102 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.debezium.connector.postgresql; + +import org.apache.flink.cdc.connectors.postgres.testutils.TestHelper; + +import io.debezium.connector.postgresql.connection.Lsn; +import io.debezium.connector.postgresql.connection.WalPositionLocator; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Unit test for the WAL-position-search guard in {@link PostgresStreamingChangeEventSource}. + * + *
On an idle publication the search loop blocks forever, because it waits for a decoded WAL
+ * message while the only mechanism that would produce one on a quiet database (the heartbeat action
+ * query) runs from the main streaming loop that the search precedes. The fix — mirroring Debezium
+ * 2.7 — only enters the search when the stored offset has actually processed a position, i.e. when
+ * {@code searchingEnabled() && offsetContext.hasCompletelyProcessedPosition()}. This test pins that
+ * decision boundary for both a fresh start and a resumed offset.
+ */
+class PostgresStreamingChangeEventSourceTest {
+
+ private PostgresConnectorConfig connectorConfig;
+ private PostgresOffsetContext.Loader offsetLoader;
+
+ @BeforeEach
+ public void beforeEach() {
+ this.connectorConfig = new PostgresConnectorConfig(TestHelper.defaultConfig().build());
+ this.offsetLoader = new PostgresOffsetContext.Loader(this.connectorConfig);
+ }
+
+ /**
+ * Builds the {@link WalPositionLocator} the same way {@code execute} does for a stored offset.
+ */
+ private static WalPositionLocator walPositionFor(PostgresOffsetContext offsetContext) {
+ Lsn lsn =
+ offsetContext.lastCompletelyProcessedLsn() != null
+ ? offsetContext.lastCompletelyProcessedLsn()
+ : offsetContext.lsn();
+ return new WalPositionLocator(offsetContext.lastCommitLsn(), lsn);
+ }
+
+ @Test
+ void shouldNotSearchWalPositionOnFreshStart() {
+ // A fresh stream split start: the starting offset carries an LSN (the low watermark) but
+ // nothing has been completely processed yet.
+ final Map