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]

Reply via email to