This is an automated email from the ASF dual-hosted git repository.

JNSimba pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris-flink-connector.git


The following commit(s) were added to refs/heads/master by this push:
     new a792b21e [Fix] Retry incremental reads after visibility wait timeout 
(#697)
a792b21e is described below

commit a792b21e0156ce9582c42f87042c3840bd489b54
Author: wudi <[email protected]>
AuthorDate: Mon Sep 14 14:15:46 2026 +0800

    [Fix] Retry incremental reads after visibility wait timeout (#697)
    
    An incremental @incr query can return ERR_INCR_VISIBLE_WAIT_TIMEOUT while a 
transaction in the requested window is still becoming visible. Retry the same 
query and window only for error code 5101, with a 1–5 second capped backoff. 
The retry budget defaults to 5 minutes and can be set with 
source.binlog.visible-wait-timeout; 0s disables connector retries. Other errors 
continue to fail immediately.
    
    The S3 integration test also uses the existing pinned MinIO release from 
quay.io/minio/minio, because the Docker Hub image returns an access-denied 
response.
---
 .../doris/flink/cfg/ConfigurationOptions.java      |  1 +
 .../apache/doris/flink/cfg/DorisReadOptions.java   | 22 +++++++++-
 .../source/reader/DorisFlightValueReader.java      | 42 +++++++++++++++++-
 .../doris/flink/table/DorisConfigOptions.java      |  8 ++++
 .../source/reader/DorisFlightValueReaderTest.java  | 50 ++++++++++++++++++++++
 .../org/apache/doris/flink/source/DorisSource.java |  3 ++
 .../flink/table/DorisDynamicTableFactory.java      |  4 ++
 .../doris/flink/source/DorisSourceOptionsTest.java | 10 +++++
 .../flink/table/DorisDynamicTableFactoryTest.java  |  4 +-
 .../org/apache/doris/flink/source/DorisSource.java |  3 ++
 .../flink/table/DorisDynamicTableFactory.java      |  4 ++
 .../flink/table/DorisDynamicTableFactoryTest.java  |  4 +-
 .../apache/doris/flink/sink/S3TvfSinkITCase.java   |  2 +-
 13 files changed, 152 insertions(+), 5 deletions(-)

diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/ConfigurationOptions.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/ConfigurationOptions.java
index ca1976fe..5a77240a 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/ConfigurationOptions.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/ConfigurationOptions.java
@@ -64,6 +64,7 @@ public interface ConfigurationOptions {
 
     String FLIGHT_SQL_PORT = "source.flight-sql-port";
     Integer FLIGHT_SQL_PORT_DEFAULT = -1;
+    Long SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT_MS_DEFAULT = 5 * 60 * 1000L;
 
     String SINK_HTTP_UTF8_CHARSET = "sink.http-utf8-charset";
     Boolean SINK_HTTP_UTF8_CHARSET_DEFAULT = false;
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisReadOptions.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisReadOptions.java
index aee0c6fc..6d03e73d 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisReadOptions.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisReadOptions.java
@@ -49,6 +49,7 @@ public class DorisReadOptions implements Serializable {
     private String scanTimestamp;
     private DorisBinlogIncrementType binlogIncrementType;
     private Long binlogPollIntervalMs;
+    private Long binlogVisibleWaitTimeoutMs;
     private String binlogOffsetTable;
     private String binlogConsumerId;
 
@@ -72,7 +73,8 @@ public class DorisReadOptions implements Serializable {
             DorisSourceScanMode scanMode,
             String scanTimestamp,
             DorisBinlogIncrementType binlogIncrementType,
-            Long binlogPollIntervalMs) {
+            Long binlogPollIntervalMs,
+            Long binlogVisibleWaitTimeoutMs) {
         this(
                 readFields,
                 filterQuery,
@@ -94,6 +96,7 @@ public class DorisReadOptions implements Serializable {
                 scanTimestamp,
                 binlogIncrementType,
                 binlogPollIntervalMs,
+                binlogVisibleWaitTimeoutMs,
                 null,
                 null);
     }
@@ -119,6 +122,7 @@ public class DorisReadOptions implements Serializable {
             String scanTimestamp,
             DorisBinlogIncrementType binlogIncrementType,
             Long binlogPollIntervalMs,
+            Long binlogVisibleWaitTimeoutMs,
             String binlogOffsetTable,
             String binlogConsumerId) {
         this.readFields = readFields;
@@ -141,6 +145,7 @@ public class DorisReadOptions implements Serializable {
         this.scanTimestamp = scanTimestamp;
         this.binlogIncrementType = binlogIncrementType;
         this.binlogPollIntervalMs = binlogPollIntervalMs;
+        this.binlogVisibleWaitTimeoutMs = binlogVisibleWaitTimeoutMs;
         this.binlogOffsetTable = binlogOffsetTable;
         this.binlogConsumerId = binlogConsumerId;
     }
@@ -241,6 +246,10 @@ public class DorisReadOptions implements Serializable {
         return binlogPollIntervalMs;
     }
 
+    public Long getBinlogVisibleWaitTimeoutMs() {
+        return binlogVisibleWaitTimeoutMs;
+    }
+
     public String getBinlogOffsetTable() {
         return binlogOffsetTable;
     }
@@ -286,6 +295,7 @@ public class DorisReadOptions implements Serializable {
                 && Objects.equals(scanTimestamp, that.scanTimestamp)
                 && binlogIncrementType == that.binlogIncrementType
                 && Objects.equals(binlogPollIntervalMs, 
that.binlogPollIntervalMs)
+                && Objects.equals(binlogVisibleWaitTimeoutMs, 
that.binlogVisibleWaitTimeoutMs)
                 && Objects.equals(binlogOffsetTable, that.binlogOffsetTable)
                 && Objects.equals(binlogConsumerId, that.binlogConsumerId);
     }
@@ -313,6 +323,7 @@ public class DorisReadOptions implements Serializable {
                 scanTimestamp,
                 binlogIncrementType,
                 binlogPollIntervalMs,
+                binlogVisibleWaitTimeoutMs,
                 binlogOffsetTable,
                 binlogConsumerId);
     }
@@ -339,6 +350,7 @@ public class DorisReadOptions implements Serializable {
                 scanTimestamp,
                 binlogIncrementType,
                 binlogPollIntervalMs,
+                binlogVisibleWaitTimeoutMs,
                 binlogOffsetTable,
                 binlogConsumerId);
     }
@@ -372,6 +384,8 @@ public class DorisReadOptions implements Serializable {
         private String scanTimestamp;
         private DorisBinlogIncrementType binlogIncrementType = 
DorisBinlogIncrementType.DETAIL;
         private Long binlogPollIntervalMs = 10_000L;
+        private Long binlogVisibleWaitTimeoutMs =
+                
ConfigurationOptions.SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT_MS_DEFAULT;
         private String binlogOffsetTable;
         private String binlogConsumerId;
 
@@ -564,6 +578,11 @@ public class DorisReadOptions implements Serializable {
             return this;
         }
 
+        public Builder setBinlogVisibleWaitTimeoutMs(Long 
binlogVisibleWaitTimeoutMs) {
+            this.binlogVisibleWaitTimeoutMs = binlogVisibleWaitTimeoutMs;
+            return this;
+        }
+
         public Builder setBinlogOffsetTable(String binlogOffsetTable) {
             this.binlogOffsetTable = binlogOffsetTable;
             return this;
@@ -601,6 +620,7 @@ public class DorisReadOptions implements Serializable {
                     scanTimestamp,
                     binlogIncrementType,
                     binlogPollIntervalMs,
+                    binlogVisibleWaitTimeoutMs,
                     binlogOffsetTable,
                     binlogConsumerId);
         }
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/reader/DorisFlightValueReader.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/reader/DorisFlightValueReader.java
index 63abef14..2dd0f2dc 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/reader/DorisFlightValueReader.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/reader/DorisFlightValueReader.java
@@ -65,6 +65,8 @@ import static 
org.apache.doris.flink.util.ErrorMessages.SHOULD_NOT_HAPPEN_MESSAG
 public class DorisFlightValueReader extends ValueReader implements 
AutoCloseable {
     private static final Logger LOG = 
LoggerFactory.getLogger(DorisFlightValueReader.class);
     private static final String PREFIX = "/* ApplicationName=Flink 
ArrowFlightSQL Query */";
+    private static final String DORIS_ERROR_CODE = "doris-error-code";
+    private static final String INCR_VISIBLE_WAIT_TIMEOUT = "5101";
 
     protected AdbcConnection client;
     private RootAllocator allocator;
@@ -121,8 +123,15 @@ public class DorisFlightValueReader extends ValueReader 
implements AutoCloseable
             } else {
                 throw new DorisRuntimeException("Unknown Doris split type: " + 
split);
             }
-            this.queryResult = statement.executeQuery();
+            this.queryResult =
+                    split instanceof DorisStreamSplit
+                            ? executeQueryWithRetry(
+                                    statement, 
readOptions.getBinlogVisibleWaitTimeoutMs())
+                            : statement.executeQuery();
             this.arrowReader = queryResult.getReader();
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            throw new RuntimeException("Interrupted while retrying Doris 
incremental query", e);
         } catch (AdbcException e) {
             throw new RuntimeException(e);
         } finally {
@@ -131,6 +140,37 @@ public class DorisFlightValueReader extends ValueReader 
implements AutoCloseable
         LOG.debug("Open scan result is, schema: {}.", schema);
     }
 
+    static AdbcStatement.QueryResult executeQueryWithRetry(
+            AdbcStatement statement, long retryTimeoutMs)
+            throws AdbcException, InterruptedException {
+        long startTimeMs = System.currentTimeMillis();
+        int retries = 1;
+        while (true) {
+            try {
+                return statement.executeQuery();
+            } catch (AdbcException e) {
+                long elapsedMs = System.currentTimeMillis() - startTimeMs;
+                if (!isVisibleWaitTimeout(e) || elapsedMs >= retryTimeoutMs) {
+                    throw e;
+                }
+                LOG.warn(
+                        "Doris incremental query visibility wait timed out; 
retrying, retry: {}, elapsed: {} ms, error: {}",
+                        retries,
+                        elapsedMs,
+                        e.getMessage());
+                Thread.sleep(Math.min(retries++, 5) * 1000L);
+            }
+        }
+    }
+
+    private static boolean isVisibleWaitTimeout(AdbcException error) {
+        return error.getDetails().stream()
+                .anyMatch(
+                        detail ->
+                                DORIS_ERROR_CODE.equals(detail.getKey())
+                                        && 
INCR_VISIBLE_WAIT_TIMEOUT.equals(detail.getValue()));
+    }
+
     private void initSchema() {
         try {
             this.schema = RestService.getSchema(options, readOptions, LOG);
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/table/DorisConfigOptions.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/table/DorisConfigOptions.java
index 7e015501..e83df655 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/table/DorisConfigOptions.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/table/DorisConfigOptions.java
@@ -43,6 +43,7 @@ import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_REQUEST_READ
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_REQUEST_RETRIES_DEFAULT;
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_TABLET_SIZE_DEFAULT;
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_THRIFT_MAX_MESSAGE_SIZE_DEFAULT;
+import static 
org.apache.doris.flink.cfg.ConfigurationOptions.SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT_MS_DEFAULT;
 import static org.apache.doris.flink.sink.writer.LoadConstants.FORMAT_KEY;
 import static org.apache.doris.flink.sink.writer.LoadConstants.JSON;
 import static 
org.apache.doris.flink.sink.writer.LoadConstants.READ_JSON_BY_LINE;
@@ -417,6 +418,13 @@ public class DorisConfigOptions {
                     .withDescription(
                             "Interval between attempts to discover the next 
Stream split; must be "
                                     + "at least 1 second");
+    public static final ConfigOption<Duration> 
SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT =
+            ConfigOptions.key("source.binlog.visible-wait-timeout")
+                    .durationType()
+                    
.defaultValue(Duration.ofMillis(SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT_MS_DEFAULT))
+                    .withDescription(
+                            "Maximum time to retry an incremental query after 
Doris reports a "
+                                    + "visible wait timeout; 0s disables 
retries");
     public static final ConfigOption<String> SOURCE_BINLOG_OFFSET_TABLE =
             ConfigOptions.key("source.binlog.offset-table")
                     .stringType()
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/reader/DorisFlightValueReaderTest.java
 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/reader/DorisFlightValueReaderTest.java
index 6cbdb76f..bdf39a2d 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/reader/DorisFlightValueReaderTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/reader/DorisFlightValueReaderTest.java
@@ -20,6 +20,10 @@ package org.apache.doris.flink.source.reader;
 import org.apache.arrow.adbc.core.AdbcConnection;
 import org.apache.arrow.adbc.core.AdbcDatabase;
 import org.apache.arrow.adbc.core.AdbcDriver;
+import org.apache.arrow.adbc.core.AdbcException;
+import org.apache.arrow.adbc.core.AdbcStatement;
+import org.apache.arrow.adbc.core.AdbcStatusCode;
+import org.apache.arrow.adbc.core.ErrorDetail;
 import org.apache.arrow.adbc.driver.flightsql.FlightSqlConnectionProperties;
 import org.apache.arrow.adbc.driver.flightsql.FlightSqlDriver;
 import org.apache.arrow.flight.Location;
@@ -41,6 +45,7 @@ import org.mockito.MockedStatic;
 import java.io.ByteArrayInputStream;
 import java.io.InputStream;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.LinkedHashSet;
 import java.util.Map;
 
@@ -53,11 +58,45 @@ import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 class DorisFlightValueReaderTest {
 
+    @Test
+    void retriesVisibleWaitTimeoutUntilQuerySucceeds() throws Exception {
+        AdbcStatement statement = mock(AdbcStatement.class);
+        AdbcStatement.QueryResult result = 
mock(AdbcStatement.QueryResult.class);
+        
when(statement.executeQuery()).thenThrow(windowError(5101)).thenReturn(result);
+
+        assertThat(DorisFlightValueReader.executeQueryWithRetry(statement, 
30_000L))
+                .isSameAs(result);
+        verify(statement, times(2)).executeQuery();
+    }
+
+    @Test
+    void doesNotRetryOtherWindowErrors() throws Exception {
+        AdbcStatement statement = mock(AdbcStatement.class);
+        AdbcException error = windowError(5100);
+        when(statement.executeQuery()).thenThrow(error);
+
+        assertThatThrownBy(() -> 
DorisFlightValueReader.executeQueryWithRetry(statement, 30_000L))
+                .isSameAs(error);
+        verify(statement).executeQuery();
+    }
+
+    @Test
+    void stopsRetryingVisibleWaitTimeoutWhenBudgetIsExhausted() throws 
Exception {
+        AdbcStatement statement = mock(AdbcStatement.class);
+        AdbcException error = windowError(5101);
+        when(statement.executeQuery()).thenThrow(error);
+
+        assertThatThrownBy(() -> 
DorisFlightValueReader.executeQueryWithRetry(statement, 0L))
+                .isSameAs(error);
+        verify(statement).executeQuery();
+    }
+
     @Test
     void closesAcquiredResourcesWhenInitializationFails() throws Exception {
         DorisSnapshotSplit split = mock(DorisSnapshotSplit.class);
@@ -224,4 +263,15 @@ class DorisFlightValueReaderTest {
                 .setTlsOptions(tlsOptions)
                 .build();
     }
+
+    private AdbcException windowError(int code) {
+        return new AdbcException(
+                "visible wait timed out",
+                null,
+                AdbcStatusCode.IO,
+                null,
+                0,
+                Collections.singletonList(
+                        new ErrorDetail("doris-error-code", 
Integer.toString(code))));
+    }
 }
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/source/DorisSource.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/source/DorisSource.java
index 1ce5ef22..49c3625f 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/source/DorisSource.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/source/DorisSource.java
@@ -230,6 +230,9 @@ public class DorisSource<OUT>
             Preconditions.checkArgument(
                     readOptions.getBinlogPollIntervalMs() >= 
MIN_BINLOG_POLL_INTERVAL_MS,
                     "source.binlog.poll-interval must be at least 1 second");
+            Preconditions.checkArgument(
+                    readOptions.getBinlogVisibleWaitTimeoutMs() >= 0,
+                    "source.binlog.visible-wait-timeout must not be negative");
             String offsetTable = readOptions.getBinlogOffsetTable();
             String consumerId = readOptions.getBinlogConsumerId();
             Preconditions.checkArgument(
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
index 8cf6ae1b..34723646 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
@@ -99,6 +99,7 @@ import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_CONS
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_INCREMENT_TYPE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_OFFSET_TABLE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_POLL_INTERVAL;
+import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT;
 import static org.apache.doris.flink.table.DorisConfigOptions.SOURCE_SCAN_MODE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_SCAN_TIMESTAMP;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_USE_OLD_API;
@@ -188,6 +189,7 @@ public final class DorisDynamicTableFactory
         options.add(SOURCE_SCAN_TIMESTAMP);
         options.add(SOURCE_BINLOG_INCREMENT_TYPE);
         options.add(SOURCE_BINLOG_POLL_INTERVAL);
+        options.add(SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT);
         options.add(SOURCE_BINLOG_OFFSET_TABLE);
         options.add(SOURCE_BINLOG_CONSUMER_ID);
         options.add(SINK_WRITE_MODE);
@@ -274,6 +276,8 @@ public final class DorisDynamicTableFactory
                         DorisBinlogIncrementType.fromOption(
                                 
readableConfig.get(SOURCE_BINLOG_INCREMENT_TYPE)))
                 
.setBinlogPollIntervalMs(readableConfig.get(SOURCE_BINLOG_POLL_INTERVAL).toMillis())
+                .setBinlogVisibleWaitTimeoutMs(
+                        
readableConfig.get(SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT).toMillis())
                 .setBinlogOffsetTable(
                         
readableConfig.getOptional(SOURCE_BINLOG_OFFSET_TABLE).orElse(null))
                 .setBinlogConsumerId(
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/source/DorisSourceOptionsTest.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/source/DorisSourceOptionsTest.java
index 5c9eb34b..df5b5701 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/source/DorisSourceOptionsTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/source/DorisSourceOptionsTest.java
@@ -82,6 +82,7 @@ class DorisSourceOptionsTest {
         assertThat(options.getScanTimestamp()).isNull();
         
assertThat(options.getBinlogIncrementType()).isEqualTo(DorisBinlogIncrementType.DETAIL);
         assertThat(options.getBinlogPollIntervalMs()).isEqualTo(10_000L);
+        assertThat(options.getBinlogVisibleWaitTimeoutMs()).isEqualTo(5 * 60 * 
1000L);
         assertThat(options.getBinlogOffsetTable()).isNull();
         assertThat(options.getBinlogConsumerId()).isNull();
     }
@@ -94,6 +95,7 @@ class DorisSourceOptionsTest {
                         .setScanTimestamp("2026-07-20 10:00:00")
                         
.setBinlogIncrementType(DorisBinlogIncrementType.MIN_DELTA)
                         .setBinlogPollIntervalMs(3_000L)
+                        .setBinlogVisibleWaitTimeoutMs(60_000L)
                         .setBinlogOffsetTable("ops.flink_source_offsets")
                         .setBinlogConsumerId("prod.sales.orders")
                         .build();
@@ -105,6 +107,7 @@ class DorisSourceOptionsTest {
         assertThat(copy.getScanTimestamp()).isEqualTo("2026-07-20 10:00:00");
         
assertThat(copy.getBinlogIncrementType()).isEqualTo(DorisBinlogIncrementType.MIN_DELTA);
         assertThat(copy.getBinlogPollIntervalMs()).isEqualTo(3_000L);
+        assertThat(copy.getBinlogVisibleWaitTimeoutMs()).isEqualTo(60_000L);
         
assertThat(copy.getBinlogOffsetTable()).isEqualTo("ops.flink_source_offsets");
         assertThat(copy.getBinlogConsumerId()).isEqualTo("prod.sales.orders");
     }
@@ -160,6 +163,13 @@ class DorisSourceOptionsTest {
                 .hasMessageContaining("at least 1 second");
         
assertThat(buildSource(DorisReadOptions.builder().setBinlogPollIntervalMs(1_000L).build()))
                 .isNotNull();
+        assertThatThrownBy(
+                        () ->
+                                buildSource(
+                                        DorisReadOptions.builder()
+                                                
.setBinlogVisibleWaitTimeoutMs(-1L)
+                                                .build()))
+                .hasMessageContaining("visible-wait-timeout must not be 
negative");
         assertThat(
                         buildSource(
                                 DorisReadOptions.builder()
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
index 01e62307..95b0e54e 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
@@ -100,6 +100,7 @@ public class DorisDynamicTableFactoryTest {
         properties.put("source.scan.timestamp", "2026-07-20 10:00:00");
         properties.put("source.binlog.increment-type", "min_delta");
         properties.put("source.binlog.poll-interval", "3s");
+        properties.put("source.binlog.visible-wait-timeout", "7m");
         DynamicTableSource actual = FactoryMocks.createTableSource(SCHEMA, 
properties);
         DorisOptions options =
                 DorisOptions.builder()
@@ -140,7 +141,8 @@ public class DorisDynamicTableFactoryTest {
                 .setScanMode(DorisSourceScanMode.FROM_TIMESTAMP)
                 .setScanTimestamp("2026-07-20 10:00:00")
                 .setBinlogIncrementType(DorisBinlogIncrementType.MIN_DELTA)
-                .setBinlogPollIntervalMs(3_000L);
+                .setBinlogPollIntervalMs(3_000L)
+                .setBinlogVisibleWaitTimeoutMs(7 * 60_000L);
         DorisDynamicTableSource expected =
                 new DorisDynamicTableSource(
                         options,
diff --git 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/source/DorisSource.java
 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/source/DorisSource.java
index 6a323a57..5c214ea1 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/source/DorisSource.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/source/DorisSource.java
@@ -225,6 +225,9 @@ public class DorisSource<OUT>
             Preconditions.checkArgument(
                     readOptions.getBinlogPollIntervalMs() >= 
MIN_BINLOG_POLL_INTERVAL_MS,
                     "source.binlog.poll-interval must be at least 1 second");
+            Preconditions.checkArgument(
+                    readOptions.getBinlogVisibleWaitTimeoutMs() >= 0,
+                    "source.binlog.visible-wait-timeout must not be negative");
             String offsetTable = readOptions.getBinlogOffsetTable();
             String consumerId = readOptions.getBinlogConsumerId();
             Preconditions.checkArgument(
diff --git 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
index ee89ecbc..72461062 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
@@ -99,6 +99,7 @@ import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_CONS
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_INCREMENT_TYPE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_OFFSET_TABLE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_POLL_INTERVAL;
+import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT;
 import static org.apache.doris.flink.table.DorisConfigOptions.SOURCE_SCAN_MODE;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_SCAN_TIMESTAMP;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SOURCE_USE_OLD_API;
@@ -188,6 +189,7 @@ public final class DorisDynamicTableFactory
         options.add(SOURCE_SCAN_TIMESTAMP);
         options.add(SOURCE_BINLOG_INCREMENT_TYPE);
         options.add(SOURCE_BINLOG_POLL_INTERVAL);
+        options.add(SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT);
         options.add(SOURCE_BINLOG_OFFSET_TABLE);
         options.add(SOURCE_BINLOG_CONSUMER_ID);
         options.add(SINK_WRITE_MODE);
@@ -274,6 +276,8 @@ public final class DorisDynamicTableFactory
                         DorisBinlogIncrementType.fromOption(
                                 
readableConfig.get(SOURCE_BINLOG_INCREMENT_TYPE)))
                 
.setBinlogPollIntervalMs(readableConfig.get(SOURCE_BINLOG_POLL_INTERVAL).toMillis())
+                .setBinlogVisibleWaitTimeoutMs(
+                        
readableConfig.get(SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT).toMillis())
                 .setBinlogOffsetTable(
                         
readableConfig.getOptional(SOURCE_BINLOG_OFFSET_TABLE).orElse(null))
                 .setBinlogConsumerId(
diff --git 
a/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
 
b/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
index b22ffc88..754c33fe 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
@@ -100,6 +100,7 @@ public class DorisDynamicTableFactoryTest {
         properties.put("source.scan.timestamp", "2026-07-20 10:00:00");
         properties.put("source.binlog.increment-type", "min_delta");
         properties.put("source.binlog.poll-interval", "3s");
+        properties.put("source.binlog.visible-wait-timeout", "7m");
         DynamicTableSource actual = FactoryMocks.createTableSource(SCHEMA, 
properties);
         DorisOptions options =
                 DorisOptions.builder()
@@ -140,7 +141,8 @@ public class DorisDynamicTableFactoryTest {
                 .setScanMode(DorisSourceScanMode.FROM_TIMESTAMP)
                 .setScanTimestamp("2026-07-20 10:00:00")
                 .setBinlogIncrementType(DorisBinlogIncrementType.MIN_DELTA)
-                .setBinlogPollIntervalMs(3_000L);
+                .setBinlogPollIntervalMs(3_000L)
+                .setBinlogVisibleWaitTimeoutMs(7 * 60_000L);
         DorisDynamicTableSource expected =
                 new DorisDynamicTableSource(
                         options,
diff --git 
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
index a297e51b..97de848b 100644
--- 
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
+++ 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
@@ -94,7 +94,7 @@ public class S3TvfSinkITCase extends AbstractITCaseService {
     private static final String ACCESS_KEY = "minioadmin";
     private static final String SECRET_KEY = "minioadmin";
     private static final int MINIO_PORT = 9000;
-    private static final String MINIO_IMAGE = 
"minio/minio:RELEASE.2024-10-13T13-34-11Z";
+    private static final String MINIO_IMAGE = 
"quay.io/minio/minio:RELEASE.2024-10-13T13-34-11Z";
 
     private static GenericContainer<?> minio;
     private static S3Client s3Client;


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to