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 9c3bb1b6 [Fix] Improve CDC TSO handling and S3 TVF upload efficiency
(#694)
9c3bb1b6 is described below
commit 9c3bb1b696500707b3368192bc0df00c1c813ecf
Author: wudi <[email protected]>
AuthorDate: Fri Sep 11 14:12:54 2026 +0800
[Fix] Improve CDC TSO handling and S3 TVF upload efficiency (#694)
Query the formatted binlog start timestamp directly from
information_schema.tso_status through the existing statement API.
Preserve the existing frontend retry and request timeout behavior, and
include complete Doris response content in statement failures and non-200 HTTP
errors.
Let the MOW E2E sink use the default Unique Key behavior that disables 2PC.
Align the S3 TVF upload queue with sink.flush.queue-size instead of using a
fixed queue length.
Upload TVF byte arrays through a repeatable content provider to avoid the
additional request-body copy while retaining retry support.
Log the object name, key, size, and upload duration for each TVF file, plus
the label, object count, attempt, and duration for each INSERT.
---
.../doris/flink/cfg/DorisExecutionOptions.java | 2 +-
.../org/apache/doris/flink/rest/RestService.java | 78 +++----------
.../flink/sink/writer/tvf/S3ClientObjectStore.java | 8 +-
.../flink/sink/writer/tvf/S3TvfCommitter.java | 16 ++-
.../doris/flink/sink/writer/tvf/S3TvfWriter.java | 35 +++++-
.../doris/flink/source/split/DorisStreamSplit.java | 4 +-
.../doris/flink/table/DorisConfigOptions.java | 5 +-
.../doris/flink/rest/DorisTsoResponseTest.java | 128 ++++++++++++---------
.../sink/writer/tvf/S3ClientObjectStoreTest.java | 11 +-
.../flink/sink/writer/tvf/S3TvfCommitterTest.java | 26 +++++
.../flink/sink/writer/tvf/S3TvfWriterTest.java | 66 ++++++++++-
.../split/DorisSourceSplitSerializerTest.java | 3 +
.../flink/sink/writer/tvf/S3TvfWriterAdapter.java | 3 +-
.../flink/sink/writer/tvf/S3TvfWriterAdapter.java | 3 +-
.../e2e/DorisIncrementalSourceE2ECase.java | 1 -
15 files changed, 250 insertions(+), 139 deletions(-)
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisExecutionOptions.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisExecutionOptions.java
index 5f7fe889..138e7cae 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisExecutionOptions.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/DorisExecutionOptions.java
@@ -496,7 +496,7 @@ public class DorisExecutionOptions implements Serializable {
}
/**
- * Set queue size in batch mode.
+ * Set queue size for asynchronous batch flush or TVF upload.
*
* @param flushQueueSize
* @return this DorisExecutionOptions.builder.
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/rest/RestService.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/rest/RestService.java
index 945c8e16..f93f2899 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/rest/RestService.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/rest/RestService.java
@@ -100,13 +100,11 @@ public class RestService implements Serializable {
private static final String QUERY_PLAN_API = "/api/%s/%s/_query_plan";
private static final String STATEMENT_EXEC_API =
"/api/query/default_cluster/information_schema";
- private static final String CURRENT_TSO_API = "/api/tso";
+ private static final String CURRENT_TIMESTAMP_SQL =
+ "SELECT FROM_UNIXTIME(CURRENT_TSO_PHYSICAL_TIME / 1000, "
+ + "'%Y-%m-%d %H:%i:%s') FROM
information_schema.tso_status";
- /**
- * Resolves the current Doris TSO to the timestamp format accepted by
row-binlog queries. 1.
- * Request the current TSO from a configured FE. 2. Call {@code
FROM_UNIXTIME} to convert it to
- * {@code yyyy-MM-dd HH:mm:ss}.
- */
+ /** Resolves the current Doris TSO to the timestamp format accepted by
row-binlog queries. */
public static String resolveCurrentTimestamp(
DorisOptions options, DorisReadOptions readOptions, Logger logger)
{
List<String> endpoints = allEndpoints(options.getFenodes(), logger);
@@ -120,17 +118,13 @@ public class RestService implements Serializable {
for (int attempt = 0; attempt < maxAttempts; attempt++) {
String endpoint = endpoints.get(attempt % endpoints.size());
try {
- long physicalTime =
- requestCurrentTsoPhysicalTime(options, readOptions,
endpoint, logger);
- // Currently, the TSO API does not return formatted time,
- // so an additional formatting step is required.
String timestamp =
parseScalarStatementResult(
executeStatementAtEndpoint(
options,
readOptions,
endpoint,
- buildCurrentTimestampSql(physicalTime),
+ CURRENT_TIMESTAMP_SQL,
logger));
return validateCurrentTimestamp(timestamp);
} catch (RuntimeException e) {
@@ -148,24 +142,6 @@ public class RestService implements Serializable {
lastFailure);
}
- private static long requestCurrentTsoPhysicalTime(
- DorisOptions options, DorisReadOptions readOptions, String
endpoint, Logger logger) {
- HttpGet request =
- new HttpGet(
- DorisUrlBuilder.buildHttpUrl(
- options.getTlsOptions(), endpoint,
CURRENT_TSO_API));
- request.setHeader(HttpHeaders.AUTHORIZATION, authHeader(options));
- request.setConfig(createRequestConfig(readOptions));
- try {
- return parseCurrentTsoPhysicalTime(
- handleResponse(request, options.getTlsOptions(),
logger).toString());
- } catch (RuntimeException e) {
- throw new DorisRuntimeException(
- "Failed to get current TSO from Doris FE " + endpoint + ":
" + e.getMessage(),
- e);
- }
- }
-
private static RequestConfig createRequestConfig(DorisReadOptions
readOptions) {
int connectTimeout =
readOptions.getRequestConnectTimeoutMs() == null
@@ -182,12 +158,6 @@ public class RestService implements Serializable {
.build();
}
- @VisibleForTesting
- static String buildCurrentTimestampSql(long physicalTime) {
- return String.format(
- "SELECT FROM_UNIXTIME(%d / 1000, '%%Y-%%m-%%d %%H:%%i:%%s')",
physicalTime);
- }
-
@VisibleForTesting
static String validateCurrentTimestamp(String timestamp) {
if (!DorisStreamSplit.isValidTimestamp(timestamp)) {
@@ -197,31 +167,6 @@ public class RestService implements Serializable {
return timestamp;
}
- @VisibleForTesting
- public static long parseCurrentTsoPhysicalTime(String response) {
- try {
- JsonNode root = objectMapper.readTree(response);
- int code = root.path("code").asInt(Integer.MIN_VALUE);
- if (code != REST_RESPONSE_CODE_OK) {
- throw new DorisRuntimeException(
- "Failed to get current Doris TSO: " +
root.path("msg").asText());
- }
- JsonNode physicalTimeNode =
root.path("data").path("current_tso_physical_time");
- if (physicalTimeNode.isMissingNode() || physicalTimeNode.isNull())
{
- throw new DorisRuntimeException(
- "Missing current_tso_physical_time in TSO response");
- }
- long physicalTime = physicalTimeNode.asLong(-1L);
- if (physicalTime <= 0) {
- throw new DorisRuntimeException(
- "Invalid current_tso_physical_time: " +
physicalTimeNode.asText());
- }
- return physicalTime;
- } catch (JsonProcessingException e) {
- throw new DorisRuntimeException("Invalid Doris TSO response", e);
- }
- }
-
/**
* send request to Doris FE and get response json string.
*
@@ -602,15 +547,20 @@ public class RestService implements Serializable {
CloseableHttpResponse response = httpclient.execute(request)) {
final int statusCode = response.getStatusLine().getStatusCode();
final String reasonPhrase =
response.getStatusLine().getReasonPhrase();
- if (statusCode == 200 && response.getEntity() != null) {
- String responseEntity =
EntityUtils.toString(response.getEntity());
+ String responseEntity =
+ response.getEntity() == null
+ ? null
+ : EntityUtils.toString(response.getEntity());
+ if (statusCode == 200 && responseEntity != null) {
return objectMapper.readTree(responseEntity);
} else {
throw new DorisRuntimeException(
"Failed to parse response, status: "
+ statusCode
+ ", reason: "
- + reasonPhrase);
+ + reasonPhrase
+ + ", response: "
+ + responseEntity);
}
} catch (Exception e) {
logger.trace("request error,", e);
@@ -653,7 +603,7 @@ public class RestService implements Serializable {
JsonNode response = handleResponse(httpPost,
options.getTlsOptions(), logger);
if (response.has("code") && response.path("code").asInt() !=
REST_RESPONSE_CODE_OK) {
throw new DorisRuntimeException(
- "Failed to execute Doris statement: " +
response.path("msg").asText());
+ "Failed to execute Doris statement, response: " +
response);
}
return response;
} catch (DorisRuntimeException e) {
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStore.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStore.java
index 4c819401..f8e25ff7 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStore.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStore.java
@@ -27,6 +27,7 @@ import software.amazon.awssdk.services.s3.S3Client;
import software.amazon.awssdk.services.s3.S3Configuration;
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
+import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.net.URI;
@@ -72,7 +73,12 @@ public class S3ClientObjectStore implements S3ObjectStore {
.contentType(JSON_LINES_CONTENT_TYPE)
.build();
try {
- s3Client.putObject(request, RequestBody.fromBytes(content));
+ s3Client.putObject(
+ request,
+ RequestBody.fromContentProvider(
+ () -> new ByteArrayInputStream(content),
+ content.length,
+ JSON_LINES_CONTENT_TYPE));
} catch (RuntimeException e) {
throw new IOException(
String.format(
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitter.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitter.java
index 85e8fbb7..4fcaf0ec 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitter.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitter.java
@@ -32,6 +32,7 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
+import java.util.concurrent.TimeUnit;
/** Commits staged objects with one INSERT statement per writer and
checkpoint. */
public class S3TvfCommitter implements Committer<S3TvfCommittable> {
@@ -88,16 +89,25 @@ public class S3TvfCommitter implements
Committer<S3TvfCommittable> {
private boolean commitOne(S3TvfCommittable committable) throws
IOException, SQLException {
String insertSql = sqlBuilder.buildInsertSql(committable);
for (int attempt = 0; attempt <= maxRetries; attempt++) {
+ long insertStartedAtNanos = System.nanoTime();
try {
loadClient.executeInsert(insertSql, sessionVariables);
- LOG.info("TVF load committed with label {}.",
committable.getLabel());
+ LOG.info(
+ "TVF insert completed, label={}, objectCount={},
attempt={}, "
+ + "insertTimeMs={}.",
+ committable.getLabel(),
+ committable.getObjectKeys().size(),
+ attempt + 1,
+ TimeUnit.NANOSECONDS.toMillis(System.nanoTime() -
insertStartedAtNanos));
return false;
} catch (SQLException e) {
LOG.warn(
- "TVF insert failed for label {} on attempt {} "
- + "(SQLState={}, errorCode={}).",
+ "TVF insert failed, label={}, objectCount={},
attempt={}, "
+ + "insertTimeMs={}, SQLState={},
errorCode={}.",
committable.getLabel(),
+ committable.getObjectKeys().size(),
attempt + 1,
+ TimeUnit.NANOSECONDS.toMillis(System.nanoTime() -
insertStartedAtNanos),
e.getSQLState(),
e.getErrorCode());
if (isLabelAlreadyUsed(e, committable.getLabel())) {
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriter.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriter.java
index 8d4df2a1..a53608ba 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriter.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriter.java
@@ -23,6 +23,8 @@ import org.apache.flink.util.concurrent.ExecutorThreadFactory;
import org.apache.doris.flink.sink.writer.DorisWriterState;
import org.apache.doris.flink.sink.writer.serializer.DorisRecord;
import org.apache.doris.flink.sink.writer.serializer.DorisRecordSerializer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
@@ -34,13 +36,14 @@ import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
/** Shared writer that stages JSON Lines files in S3-compatible object
storage. */
public class S3TvfWriter<IN> {
+ private static final Logger LOG =
LoggerFactory.getLogger(S3TvfWriter.class);
private static final byte NEW_LINE = '\n';
- private static final int UPLOAD_QUEUE_SIZE = 1;
private final int subtaskId;
private final DorisRecordSerializer<IN> serializer;
@@ -52,10 +55,10 @@ public class S3TvfWriter<IN> {
private final List<String> columns;
private final boolean deleteSignEnabled;
private final int maxBytes;
+ private final int uploadQueueSize;
private final ByteArrayOutputStream buffer = new ByteArrayOutputStream();
private final List<String> currentObjectKeys = new ArrayList<>();
- private final BlockingQueue<Runnable> uploadQueue =
- new LinkedBlockingQueue<>(UPLOAD_QUEUE_SIZE);
+ private final BlockingQueue<Runnable> uploadQueue;
private final AtomicReference<IOException> uploadException = new
AtomicReference<>();
private final ExecutorService uploadExecutor;
@@ -73,8 +76,10 @@ public class S3TvfWriter<IN> {
String labelPrefix,
List<String> columns,
boolean deleteSignEnabled,
- int maxBytes) {
+ int maxBytes,
+ int uploadQueueSize) {
Preconditions.checkArgument(maxBytes > 0, "TVF buffer max bytes must
be positive.");
+ Preconditions.checkArgument(uploadQueueSize > 0, "TVF upload queue
size must be positive.");
this.currentCheckpointId = restoredCheckpointId + 1;
this.subtaskId = subtaskId;
this.serializer = serializer;
@@ -86,6 +91,8 @@ public class S3TvfWriter<IN> {
this.columns = Collections.unmodifiableList(new ArrayList<>(columns));
this.deleteSignEnabled = deleteSignEnabled;
this.maxBytes = maxBytes;
+ this.uploadQueueSize = uploadQueueSize;
+ this.uploadQueue = new LinkedBlockingQueue<>(uploadQueueSize);
this.uploadExecutor =
Executors.newSingleThreadExecutor(
new ExecutorThreadFactory("s3-tvf-upload-" +
subtaskId));
@@ -175,10 +182,28 @@ public class S3TvfWriter<IN> {
if (uploadException.get() != null) {
return;
}
+ long uploadStartedAtNanos = System.nanoTime();
try {
objectStore.put(objectKey, content);
currentObjectKeys.add(objectKey);
+ LOG.info(
+ "S3 TVF object upload completed, fileName={},
objectKey={}, "
+ + "sizeBytes={}, uploadTimeMs={}.",
+ fileName,
+ objectKey,
+ content.length,
+ TimeUnit.NANOSECONDS.toMillis(
+ System.nanoTime() -
uploadStartedAtNanos));
} catch (Exception e) {
+ LOG.warn(
+ "S3 TVF object upload failed, fileName={},
objectKey={}, "
+ + "sizeBytes={}, uploadTimeMs={}.",
+ fileName,
+ objectKey,
+ content.length,
+ TimeUnit.NANOSECONDS.toMillis(
+ System.nanoTime() -
uploadStartedAtNanos),
+ e);
IOException failure =
e instanceof IOException
? (IOException) e
@@ -201,7 +226,7 @@ public class S3TvfWriter<IN> {
}
private void waitForUploads() throws IOException {
- for (int i = 0; i <= UPLOAD_QUEUE_SIZE; i++) {
+ for (int i = 0; i <= uploadQueueSize; i++) {
putUpload(() -> {});
}
checkUploadException();
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/split/DorisStreamSplit.java
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/split/DorisStreamSplit.java
index 99fa6c4b..f4e3ca66 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/split/DorisStreamSplit.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/source/split/DorisStreamSplit.java
@@ -20,7 +20,7 @@ package org.apache.doris.flink.source.split;
import java.util.Objects;
import java.util.regex.Pattern;
-/** A finite Doris row-binlog query range with an exclusive start and
inclusive end. */
+/** A finite Doris row-binlog query range with an inclusive start and
exclusive end. */
public final class DorisStreamSplit implements DorisSourceSplit {
private static final Pattern TIMESTAMP_PATTERN =
Pattern.compile("\\d{4}-\\d{2}-\\d{2} \\d{2}:\\d{2}:\\d{2}");
@@ -102,6 +102,6 @@ public final class DorisStreamSplit implements
DorisSourceSplit {
@Override
public String toString() {
- return "DorisStreamSplit{" + splitId + ", (" + startTimestamp + ", " +
endTimestamp + "]}";
+ return "DorisStreamSplit{" + splitId + ", [" + startTimestamp + ", " +
endTimestamp + ")}";
}
}
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 c74be29e..7e015501 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
@@ -352,7 +352,8 @@ public class DorisConfigOptions {
ConfigOptions.key("sink.flush.queue-size")
.intType()
.defaultValue(2)
- .withDescription("Queue length for async stream load,
default is 2");
+ .withDescription(
+ "Queue length for asynchronous stream load or TVF
upload, default is 2");
public static final ConfigOption<Integer> SINK_BUFFER_FLUSH_MAX_ROWS =
ConfigOptions.key("sink.buffer-flush.max-rows")
@@ -403,7 +404,7 @@ public class DorisConfigOptions {
ConfigOptions.key("source.scan.timestamp")
.stringType()
.noDefaultValue()
- .withDescription("Exclusive start timestamp for
from-timestamp mode");
+ .withDescription("Inclusive start timestamp for
from-timestamp mode");
public static final ConfigOption<String> SOURCE_BINLOG_INCREMENT_TYPE =
ConfigOptions.key("source.binlog.increment-type")
.stringType()
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/rest/DorisTsoResponseTest.java
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/rest/DorisTsoResponseTest.java
index 37bb63fa..8055d46c 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/rest/DorisTsoResponseTest.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/rest/DorisTsoResponseTest.java
@@ -18,8 +18,10 @@
package org.apache.doris.flink.rest;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.sun.net.httpserver.HttpServer;
import org.apache.doris.flink.cfg.DorisOptions;
import org.apache.doris.flink.cfg.DorisReadOptions;
+import org.apache.http.client.methods.HttpGet;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.client.methods.HttpRequestBase;
import org.apache.http.util.EntityUtils;
@@ -28,10 +30,13 @@ import org.mockito.MockedStatic;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
import java.util.concurrent.atomic.AtomicInteger;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.assertj.core.api.Assertions.catchThrowable;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.CALLS_REAL_METHODS;
import static org.mockito.Mockito.mockStatic;
@@ -41,46 +46,66 @@ class DorisTsoResponseTest {
private static final Logger LOG =
LoggerFactory.getLogger(DorisTsoResponseTest.class);
@Test
- void extractsOnlyPhysicalTime() {
- String response =
- "{\"code\":0,\"msg\":\"success\",\"data\":{"
- + "\"current_tso\":461373440032243713,"
- + "\"current_tso_physical_time\":1760000000123,"
- + "\"current_tso_logical_counter\":1}}";
-
-
assertThat(RestService.parseCurrentTsoPhysicalTime(response)).isEqualTo(1760000000123L);
+ void validatesTimestampFormattingResult() {
+ assertThat(RestService.validateCurrentTimestamp("2026-07-20 10:00:00"))
+ .isEqualTo("2026-07-20 10:00:00");
+ assertThatThrownBy(() ->
RestService.validateCurrentTimestamp("2026-07-20 10:00:00.123000"))
+ .hasMessageContaining("yyyy-MM-dd HH:mm:ss");
}
@Test
- void rejectsErrorAndMissingPhysicalTime() {
- assertThatThrownBy(
- () ->
- RestService.parseCurrentTsoPhysicalTime(
- "{\"code\":1,\"msg\":\"Temporary
failure\"}"))
- .hasMessageContaining("Temporary failure");
- assertThatThrownBy(
- () ->
- RestService.parseCurrentTsoPhysicalTime(
-
"{\"code\":0,\"msg\":\"success\",\"data\":{}}"))
- .hasMessageContaining("current_tso_physical_time");
- }
+ void includesCompleteStatementResponseInError() throws Exception {
+ DorisOptions options =
+ DorisOptions.builder()
+ .setFenodes("frontend:8030")
+ .setUsername("root")
+ .setPassword("")
+ .build();
+ DorisReadOptions readOptions =
DorisReadOptions.builder().setRequestRetries(1).build();
+ String response =
+ "{\"code\":1,\"msg\":\"Error\",\"data\":\"Table [tso_status]
does not exist\"}";
- @Test
- void buildsTimestampFormattingSql() {
- assertThat(RestService.buildCurrentTimestampSql(1760000000123L))
- .isEqualTo("SELECT FROM_UNIXTIME(1760000000123 / 1000,
'%Y-%m-%d %H:%i:%s')");
+ try (MockedStatic<RestService> mocked = mockStatic(RestService.class,
CALLS_REAL_METHODS)) {
+ mocked.when(() -> RestService.handleResponse(any(), any(), any()))
+ .thenReturn(new ObjectMapper().readTree(response));
+
+ Throwable error =
+ catchThrowable(
+ () -> RestService.resolveCurrentTimestamp(options,
readOptions, LOG));
+
+ assertThat(error)
+ .hasMessage("Failed to resolve current Doris timestamp
after 1 attempts");
+ assertThat(error.getCause()).hasMessageContaining(response);
+ }
}
@Test
- void validatesTimestampFormattingResult() {
- assertThat(RestService.validateCurrentTimestamp("2026-07-20 10:00:00"))
- .isEqualTo("2026-07-20 10:00:00");
- assertThatThrownBy(() ->
RestService.validateCurrentTimestamp("2026-07-20 10:00:00.123000"))
- .hasMessageContaining("yyyy-MM-dd HH:mm:ss");
+ void includesHttpErrorResponseBody() throws Exception {
+ byte[] response =
+ "{\"code\":500,\"msg\":\"Error\",\"data\":\"Detailed
failure\"}"
+ .getBytes(StandardCharsets.UTF_8);
+ HttpServer server = HttpServer.create(new
InetSocketAddress("localhost", 0), 0);
+ server.createContext(
+ "/",
+ exchange -> {
+ exchange.sendResponseHeaders(500, response.length);
+ exchange.getResponseBody().write(response);
+ exchange.close();
+ });
+ server.start();
+
+ try {
+ HttpGet request =
+ new HttpGet("http://localhost:" +
server.getAddress().getPort() + "/");
+ assertThatThrownBy(() -> RestService.handleResponse(request, LOG))
+ .hasMessageContaining(new String(response,
StandardCharsets.UTF_8));
+ } finally {
+ server.stop(0);
+ }
}
@Test
- void retriesConfiguredFrontendAndAppliesTimeoutsAfterTsoFailure() throws
Exception {
+ void queriesTimestampFromTsoStatusWithConfiguredRetriesAndTimeouts()
throws Exception {
DorisOptions options =
DorisOptions.builder()
.setFenodes("frontend:8030")
@@ -93,8 +118,7 @@ class DorisTsoResponseTest {
.setRequestReadTimeoutMs(2345)
.setRequestRetries(2)
.build();
- AtomicInteger tsoCalls = new AtomicInteger();
- AtomicInteger formatTimestampCalls = new AtomicInteger();
+ AtomicInteger statementCalls = new AtomicInteger();
ObjectMapper mapper = new ObjectMapper();
try (MockedStatic<RestService> mocked = mockStatic(RestService.class,
CALLS_REAL_METHODS)) {
@@ -106,36 +130,30 @@ class DorisTsoResponseTest {
assertThat(request.getConfig().getConnectTimeout()).isEqualTo(1234);
assertThat(request.getConfig().getSocketTimeout()).isEqualTo(2345);
- if (request instanceof HttpPost) {
- String statement =
- EntityUtils.toString(((HttpPost)
request).getEntity());
- if (statement.contains("FROM_UNIXTIME")) {
- formatTimestampCalls.incrementAndGet();
- assertThat(request.getURI().getHost())
- .isEqualTo("frontend");
- return mapper.readTree(
-
"{\"code\":0,\"data\":{\"data\":"
- + "[[\"2026-07-20
10:00:00\"]]}}");
- }
- } else if
("/api/tso".equals(request.getURI().getPath())) {
-
assertThat(request.getURI().getHost()).isEqualTo("frontend");
- if (tsoCalls.getAndIncrement() == 0) {
- return mapper.readTree(
-
"{\"code\":1,\"msg\":\"Temporary failure\"}");
- }
+ if (!(request instanceof HttpPost)) {
+ throw new AssertionError(
+ "Expected statement request but
got "
+ + request.getURI());
+ }
+ String statement =
+ EntityUtils.toString(((HttpPost)
request).getEntity());
+ assertThat(statement)
+ .contains("FROM_UNIXTIME")
+
.contains("information_schema.tso_status");
+
assertThat(request.getURI().getHost()).isEqualTo("frontend");
+ if (statementCalls.getAndIncrement() == 0) {
return mapper.readTree(
- "{\"code\":0,\"data\":{"
- +
"\"current_tso_physical_time\":"
- + "1760000000123}}");
+ "{\"code\":1,\"msg\":\"Temporary
failure\"}");
}
- throw new AssertionError("Unexpected request:
" + request.getURI());
+ return mapper.readTree(
+ "{\"code\":0,\"data\":{\"data\":"
+ + "[[\"2026-07-20
10:00:00\"]]}}");
});
assertThat(RestService.resolveCurrentTimestamp(options,
readOptions, LOG))
.isEqualTo("2026-07-20 10:00:00");
}
- assertThat(tsoCalls.get()).isEqualTo(2);
- assertThat(formatTimestampCalls.get()).isEqualTo(1);
+ assertThat(statementCalls.get()).isEqualTo(2);
}
}
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStoreTest.java
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStoreTest.java
index 1e8e429e..a3c83f14 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStoreTest.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3ClientObjectStoreTest.java
@@ -33,7 +33,7 @@ import static org.mockito.Mockito.verify;
public class S3ClientObjectStoreTest {
@Test
- public void testPutObject() throws Exception {
+ public void testPutObjectUsesRepeatableContentProviderWithoutCopying()
throws Exception {
S3Client s3Client = mock(S3Client.class);
S3ClientObjectStore objectStore = new S3ClientObjectStore(s3Client,
"bucket");
byte[] content = "{\"id\":1}\n".getBytes(StandardCharsets.UTF_8);
@@ -47,10 +47,17 @@ public class S3ClientObjectStoreTest {
Assert.assertEquals("bucket", requestCaptor.getValue().bucket());
Assert.assertEquals("prefix_tbl_0_1_0.json",
requestCaptor.getValue().key());
Assert.assertEquals("application/x-ndjson",
requestCaptor.getValue().contentType());
- try (InputStream input =
bodyCaptor.getValue().contentStreamProvider().newStream()) {
+
+ content[0] = '[';
+ try (InputStream input =
bodyCaptor.getValue().contentStreamProvider().newStream();
+ InputStream retryInput =
+
bodyCaptor.getValue().contentStreamProvider().newStream()) {
byte[] actual = new byte[content.length];
+ byte[] retryActual = new byte[content.length];
Assert.assertEquals(content.length, input.read(actual));
+ Assert.assertEquals(content.length, retryInput.read(retryActual));
Assert.assertArrayEquals(content, actual);
+ Assert.assertArrayEquals(content, retryActual);
}
objectStore.close();
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitterTest.java
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitterTest.java
index 6b055bba..8168d8cf 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitterTest.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfCommitterTest.java
@@ -18,10 +18,13 @@
package org.apache.doris.flink.sink.writer.tvf;
import org.apache.flink.api.connector.sink2.Committer.CommitRequest;
+import org.apache.flink.testutils.logging.TestLoggerResource;
import org.apache.doris.flink.cfg.S3TvfOptions;
import org.junit.Assert;
+import org.junit.Rule;
import org.junit.Test;
+import org.slf4j.event.Level;
import java.sql.SQLException;
import java.util.ArrayDeque;
@@ -33,6 +36,29 @@ import java.util.Queue;
public class S3TvfCommitterTest {
+ @Rule
+ public final TestLoggerResource testLogger =
+ new TestLoggerResource(S3TvfCommitter.class, Level.INFO);
+
+ @Test
+ public void testLogsInsertMetrics() throws Exception {
+ RecordingLoadClient loadClient = new RecordingLoadClient();
+ S3TvfCommitter committer = createCommitter(loadClient, 3);
+
+ committer.commit(
+ Collections.singletonList(
+ request(committable("prefix_tbl_0_7_0.json",
"label_tbl_0_7"))));
+
+ Assert.assertTrue(
+ testLogger.getMessages().stream()
+ .anyMatch(
+ message ->
+ message.matches(
+ "TVF insert completed,
label=label_tbl_0_7, "
+ + "objectCount=1,
attempt=1, "
+ +
"insertTimeMs=\\d+\\.")));
+ }
+
@Test
public void testCommitsWriterRequestsIndependently() throws Exception {
RecordingLoadClient loadClient = new RecordingLoadClient();
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterTest.java
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterTest.java
index ea793c77..565e123b 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterTest.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterTest.java
@@ -17,10 +17,14 @@
package org.apache.doris.flink.sink.writer.tvf;
+import org.apache.flink.testutils.logging.TestLoggerResource;
+
import org.apache.doris.flink.sink.writer.serializer.DorisRecord;
import org.apache.doris.flink.sink.writer.serializer.DorisRecordSerializer;
import org.junit.Assert;
+import org.junit.Rule;
import org.junit.Test;
+import org.slf4j.event.Level;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
@@ -37,6 +41,29 @@ import java.util.concurrent.TimeoutException;
public class S3TvfWriterTest {
+ @Rule
+ public final TestLoggerResource testLogger =
+ new TestLoggerResource(S3TvfWriter.class, Level.INFO);
+
+ @Test
+ public void testLogsUploadedObjectMetrics() throws Exception {
+ RecordingObjectStore objectStore = new RecordingObjectStore();
+ S3TvfWriter<String> writer = createWriter(6L, objectStore);
+
+ writer.write("12345");
+ writer.flush();
+
+ Assert.assertTrue(
+ testLogger.getMessages().stream()
+ .anyMatch(
+ message ->
+ message.matches(
+ "S3 TVF object upload
completed, "
+ +
"fileName=label_tbl_2_7_0\\.json, "
+ +
"objectKey=prefix/label_tbl_2_7_0\\.json, "
+ + "sizeBytes=6,
uploadTimeMs=\\d+\\.")));
+ }
+
@Test
public void testFlushByBytesAndBuildDeterministicCommittable() throws
Exception {
RecordingObjectStore objectStore = new RecordingObjectStore();
@@ -135,8 +162,44 @@ public class S3TvfWriterTest {
}
}
+ @Test
+ public void testConfiguredUploadQueueSizeAllowsPendingUploads() throws
Exception {
+ BlockingObjectStore objectStore = new BlockingObjectStore();
+ S3TvfWriter<String> writer = createWriter(6L, objectStore, 2, 6);
+ ExecutorService caller = Executors.newSingleThreadExecutor();
+
+ try {
+ writer.write("12345");
+ Assert.assertTrue(objectStore.uploadStarted.await(5,
TimeUnit.SECONDS));
+
+ Future<?> pendingWrites =
+ caller.submit(
+ () -> {
+ writer.write("12345");
+ writer.write("12345");
+ return null;
+ });
+ pendingWrites.get(1, TimeUnit.SECONDS);
+
+ objectStore.allowUpload.countDown();
+ writer.flush();
+ } finally {
+ objectStore.allowUpload.countDown();
+ caller.shutdownNow();
+ writer.close();
+ }
+ }
+
private static S3TvfWriter<String> createWriter(
long restoredCheckpointId, RecordingObjectStore objectStore) {
+ return createWriter(restoredCheckpointId, objectStore, 2, 10);
+ }
+
+ private static S3TvfWriter<String> createWriter(
+ long restoredCheckpointId,
+ RecordingObjectStore objectStore,
+ int uploadQueueSize,
+ int maxBytes) {
DorisRecordSerializer<String> serializer =
value ->
DorisRecord.of(value.getBytes(StandardCharsets.UTF_8));
return new S3TvfWriter<>(
@@ -150,7 +213,8 @@ public class S3TvfWriterTest {
"label",
Arrays.asList("id", "name"),
true,
- 10);
+ maxBytes,
+ uploadQueueSize);
}
private static class RecordingObjectStore implements S3ObjectStore {
diff --git
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/split/DorisSourceSplitSerializerTest.java
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/split/DorisSourceSplitSerializerTest.java
index aa82c633..d78303a9 100644
---
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/split/DorisSourceSplitSerializerTest.java
+++
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/source/split/DorisSourceSplitSerializerTest.java
@@ -63,6 +63,9 @@ class DorisSourceSplitSerializerTest {
DorisStreamSplit split = DorisStreamSplit.of("2026-07-20 10:00:00",
"2026-07-20 10:00:10");
assertThat(split.splitId()).isEqualTo("stream-20260720100000-20260720100010");
+ assertThat(split.toString())
+ .isEqualTo(
+
"DorisStreamSplit{stream-20260720100000-20260720100010, [2026-07-20 10:00:00,
2026-07-20 10:00:10)}");
assertThat(roundTrip(split)).isEqualTo(split);
}
diff --git
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
index e570899d..7cb2d8b1 100644
---
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
+++
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
@@ -74,7 +74,8 @@ public class S3TvfWriterAdapter<IN>
executionOptions.getLabelPrefix(),
rowDataSerializer.getSelectedColumns(),
rowDataSerializer.isDeleteSignEnabled(),
- executionOptions.getBufferFlushMaxBytes());
+ executionOptions.getBufferFlushMaxBytes(),
+ executionOptions.getFlushQueueSize());
}
@Override
diff --git
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
index 0e13dff7..718eb147 100644
---
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
+++
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfWriterAdapter.java
@@ -74,7 +74,8 @@ public class S3TvfWriterAdapter<IN>
executionOptions.getLabelPrefix(),
rowDataSerializer.getSelectedColumns(),
rowDataSerializer.isDeleteSignEnabled(),
- executionOptions.getBufferFlushMaxBytes());
+ executionOptions.getBufferFlushMaxBytes(),
+ executionOptions.getFlushQueueSize());
}
@Override
diff --git
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/container/e2e/DorisIncrementalSourceE2ECase.java
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/container/e2e/DorisIncrementalSourceE2ECase.java
index cc0c6284..4d372081 100644
---
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/container/e2e/DorisIncrementalSourceE2ECase.java
+++
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/container/e2e/DorisIncrementalSourceE2ECase.java
@@ -195,7 +195,6 @@ public class DorisIncrementalSourceE2ECase extends
AbstractITCaseService {
+ " 'username' = '%s',\n"
+ " 'password' = '%s',\n"
+ " 'sink.label-prefix' = '%s',\n"
- + " 'sink.enable-2pc' = 'true',\n"
+ " 'sink.enable-delete' = 'true',\n"
+ " 'sink.ignore.update-before' =
'true',\n"
+ " 'sink.buffer-flush.interval' = '1s'\n"
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]