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]