From 4f2645fc863e575617bf56a94af091b6154dba8e Mon Sep 17 00:00:00 2001 From: Joao Boto Date: Wed, 2 Sep 2026 14:29:25 +0200 Subject: [PATCH] [FLINK-40538][postgres] Skip WAL position search on idle publication for a fresh start PostgresSourceFetchTaskContext.loadStartingOffsetState always returns a non-null PostgresOffsetContext built from the stream split's starting offset, so in the forked PostgresStreamingChangeEventSource#execute the WAL position search branch was taken unconditionally. On an idle publication that search loops forever waiting for a decoded message: no WAL is produced, and the heartbeat action query that would generate some only runs from the main streaming loop, which the search precedes. The job stays RUNNING while the slot's confirmed_flush_lsn never advances and WAL grows without bound. Guard the search with offsetContext.hasCompletelyProcessedPosition(), mirroring the fix Debezium shipped in 2.7. On a fresh start (nothing processed yet) the search is skipped; streaming still starts from the stored LSN, so no events are missed. On a resumed offset the search still runs. --- .../flink-connector-postgres-cdc/pom.xml | 10 + .../PostgresStreamingChangeEventSource.java | 27 ++- ...ostgresStreamingChangeEventSourceTest.java | 182 ++++++++++++++++++ 3 files changed, 214 insertions(+), 5 deletions(-) create mode 100644 flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/io/debezium/connector/postgresql/PostgresStreamingChangeEventSourceTest.java diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xml b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xml index ec73c783ac9..07fe11266e5 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xml +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xml @@ -29,6 +29,9 @@ limitations under the License. flink-connector-postgres-cdc jar + + 3.12.4 + @@ -191,6 +194,13 @@ limitations under the License. test + + org.mockito + mockito-core + ${mockito.version} + test + + 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..444caa20b11 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,8 +187,17 @@ public void execute( this.lastCompletelyProcessedLsn = replicationStream.get().startLsn(); - if (walPosition.searchingEnabled()) { - searchWalPosition(context, stream, walPosition); + // 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(); it relies on the DBZ-6635 heartbeat fix backported in + // processMessages(), so the heartbeat action query can generate WAL meanwhile. + if (walPosition.searchingEnabled() && offsetContext.hasCompletelyProcessedPosition()) { + searchWalPosition(context, partition, offsetContext, stream, walPosition); try { if (!isInPreSnapshotCatchUpStreaming(offsetContext)) { connection.commit(); @@ -356,9 +365,12 @@ private void processMessages( noMessageIterations = 0; lsnFlushingAllowed = true; } else { - if (offsetContext.hasCompletelyProcessedPosition()) { - dispatcher.dispatchHeartbeatEvent(partition, offsetContext); - } + // Dispatch heartbeats also before any WAL message was processed (backport of + // Debezium DBZ-6635, shipped in 2.4): on a fresh start over an idle publication + // this is the only path that runs the heartbeat action query, which in turn + // generates the WAL that lets streaming make progress. Without it, skipping the + // WAL-position search below would only move the stall into this loop. + dispatcher.dispatchHeartbeatEvent(partition, offsetContext); noMessageIterations++; if (noMessageIterations >= THROTTLE_NO_MESSAGE_BEFORE_PAUSE) { noMessageIterations = 0; @@ -380,6 +392,8 @@ private void processMessages( private void searchWalPosition( ChangeEventSourceContext context, + PostgresPartition partition, + PostgresOffsetContext offsetContext, final ReplicationStream stream, final WalPositionLocator walPosition) throws SQLException, InterruptedException { @@ -399,6 +413,9 @@ private void searchWalPosition( if (receivedMessage) { noMessageIterations = 0; } else { + // Dispatch heartbeats while searching so that the heartbeat action query can + // generate WAL on an idle publication and let the search terminate (Debezium 2.7). + dispatcher.dispatchHeartbeatEvent(partition, offsetContext); noMessageIterations++; if (noMessageIterations >= THROTTLE_NO_MESSAGE_BEFORE_PAUSE) { noMessageIterations = 0; 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..6a2b7032680 --- /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,182 @@ +/* + * 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.PostgresConnection; +import io.debezium.connector.postgresql.connection.PostgresReplicationConnection; +import io.debezium.connector.postgresql.connection.ReplicationStream; +import io.debezium.connector.postgresql.connection.WalPositionLocator; +import io.debezium.connector.postgresql.spi.Snapshotter; +import io.debezium.pipeline.ErrorHandler; +import io.debezium.pipeline.source.spi.ChangeEventSource; +import io.debezium.relational.TableId; +import io.debezium.util.Clock; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Unit test for the idle-publication fixes in {@link PostgresStreamingChangeEventSource}. + * + *

On an idle publication the WAL-position 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 heartbeat dispatch. The fix backports the two upstream Debezium + * changes: + * + *

+ */ +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); + } + + private static PostgresOffsetContext freshStartOffsetContext( + PostgresOffsetContext.Loader loader) { + // A fresh stream split start: the starting offset carries an LSN (the low watermark) but + // nothing has been completely processed yet. + final Map offsetValues = new HashMap<>(); + offsetValues.put(SourceInfo.LSN_KEY, 12345L); + offsetValues.put(SourceInfo.TIMESTAMP_USEC_KEY, 67890L); + return loader.load(offsetValues); + } + + @Test + void shouldNotSearchWalPositionOnFreshStart() { + final PostgresOffsetContext offsetContext = freshStartOffsetContext(offsetLoader); + + // searchingEnabled() alone is true, so the pre-fix condition would enter the search loop + // and stall forever on an idle publication... + assertThat(walPositionFor(offsetContext).searchingEnabled()).isTrue(); + // ...but the added guard is false on a fresh start, so the search is correctly skipped. + assertThat(offsetContext.hasCompletelyProcessedPosition()) + .as( + "WAL search must be skipped on a fresh start so an idle publication cannot stall it") + .isFalse(); + } + + @Test + void shouldSearchWalPositionWhenResumingFromProcessedOffset() { + // A resumed offset (e.g. after a checkpoint): a position has already been processed, so the + // search is still required to locate the exact resume point among already-seen LSNs. + final Map offsetValues = new HashMap<>(); + offsetValues.put(SourceInfo.LSN_KEY, 12345L); + offsetValues.put(SourceInfo.TIMESTAMP_USEC_KEY, 67890L); + offsetValues.put(PostgresOffsetContext.LAST_COMPLETELY_PROCESSED_LSN_KEY, 12345L); + + final PostgresOffsetContext offsetContext = offsetLoader.load(offsetValues); + + // Both the pre-fix condition and the added guard are true, so the search still runs. + assertThat(walPositionFor(offsetContext).searchingEnabled()).isTrue(); + assertThat(offsetContext.hasCompletelyProcessedPosition()) + .as("WAL search must still run when resuming from an already-processed position") + .isTrue(); + } + + @Test + void shouldStreamAndDispatchHeartbeatsOnIdlePublicationFreshStart() throws Exception { + // Fresh start over an idle publication: the WAL-position search must be skipped and the + // main streaming loop must keep dispatching heartbeats (which run the heartbeat action + // query) even though nothing has a completely processed position. Without the DBZ-6635 + // backport the heartbeat dispatch stays blocked behind hasCompletelyProcessedPosition(), + // no WAL is ever produced and the job stalls in the main loop instead of the search. + final PostgresOffsetContext offsetContext = freshStartOffsetContext(offsetLoader); + assertThat(offsetContext.hasCompletelyProcessedPosition()).isFalse(); + + final PostgresTaskContext taskContext = mock(PostgresTaskContext.class); + when(taskContext.getConfig()).thenReturn(connectorConfig); + @SuppressWarnings("unchecked") + final PostgresEventDispatcher dispatcher = mock(PostgresEventDispatcher.class); + final Snapshotter snapshotter = mock(Snapshotter.class); + when(snapshotter.shouldStream()).thenReturn(true); + final PostgresConnection connection = mock(PostgresConnection.class); + final ErrorHandler errorHandler = mock(ErrorHandler.class); + final PostgresReplicationConnection replicationConnection = + mock(PostgresReplicationConnection.class); + final ReplicationStream stream = mock(ReplicationStream.class); + // Simulate an idle publication: no message is ever received. + when(stream.readPending(any())).thenReturn(false); + when(stream.startLsn()).thenReturn(Lsn.valueOf(12345L)); + when(replicationConnection.startStreaming(any(Lsn.class), any(WalPositionLocator.class))) + .thenReturn(stream); + + final PostgresStreamingChangeEventSource source = + new PostgresStreamingChangeEventSource( + connectorConfig, + snapshotter, + connection, + dispatcher, + errorHandler, + Clock.system(), + mock(PostgresSchema.class), + taskContext, + replicationConnection); + + // Run exactly two polls of the streaming loop, then stop the source. + final AtomicInteger polls = new AtomicInteger(); + final ChangeEventSource.ChangeEventSourceContext context = + () -> polls.getAndIncrement() < 2; + + source.execute(context, new PostgresPartition("test_server"), offsetContext); + + // execute() swallows throwables into the error handler. + verify(errorHandler, never()).setProducerThrowable(any()); + // One heartbeat dispatch per poll, although no message was ever received. + verify(dispatcher, times(2)).dispatchHeartbeatEvent(any(), eq(offsetContext)); + // The WAL-position search was skipped: streaming kept the initial connection. + verify(replicationConnection, never()).reconnect(); + } +}