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-kafka-connector.git


The following commit(s) were added to refs/heads/master by this push:
     new a09a100  [Improve] Improve S3 TVF upload efficiency and observability 
(#105)
a09a100 is described below

commit a09a100de705ff0c3093ff1649759c0792ab14be
Author: wudi <[email protected]>
AuthorDate: Fri Sep 11 14:14:51 2026 +0800

    [Improve] Improve S3 TVF upload efficiency and observability (#105)
---
 .../kafka/connector/writer/AsyncS3TvfWriter.java   | 19 +++++++++++
 .../connector/writer/s3/S3ClientObjectStore.java   |  8 ++++-
 .../doris/kafka/connector/writer/s3/S3TvfLoad.java | 18 ++++++++--
 .../connector/writer/AsyncS3TvfWriterTest.java     | 32 ++++++++++++++++++
 .../writer/s3/S3ClientObjectStoreTest.java         | 23 ++++++++-----
 .../kafka/connector/writer/s3/S3TvfLoadTest.java   | 39 ++++++++++++++++++++++
 6 files changed, 127 insertions(+), 12 deletions(-)

diff --git 
a/src/main/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriter.java 
b/src/main/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriter.java
index 15e3866..19c8ca2 100644
--- 
a/src/main/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriter.java
+++ 
b/src/main/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriter.java
@@ -30,6 +30,7 @@ import java.util.concurrent.CountDownLatch;
 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;
 import org.apache.doris.kafka.connector.cfg.DorisOptions;
 import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
@@ -233,10 +234,28 @@ public class AsyncS3TvfWriter extends DorisWriter {
                     if (exception.get() != null) {
                         return;
                     }
+                    long uploadStartedAtNanos = System.nanoTime();
                     try {
                         objectStore.put(objectKey, content);
                         uploadedObjectKeys.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);
                         exception.compareAndSet(
                                 null,
                                 new DorisException("Failed to upload S3 TVF 
file " + objectKey, e));
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStore.java
 
b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStore.java
index 1c3059e..bb4aa80 100644
--- 
a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStore.java
+++ 
b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStore.java
@@ -19,6 +19,7 @@
 
 package org.apache.doris.kafka.connector.writer.s3;
 
+import java.io.ByteArrayInputStream;
 import java.io.IOException;
 import java.net.URI;
 import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
@@ -56,7 +57,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("Failed to upload staged S3 object " + 
objectKey, e);
         }
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoad.java 
b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoad.java
index 80d18aa..1f3942e 100644
--- a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoad.java
+++ b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoad.java
@@ -31,6 +31,7 @@ import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Locale;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 import java.util.regex.Pattern;
 import org.apache.doris.kafka.connector.cfg.DorisOptions;
 import org.apache.doris.kafka.connector.connection.ConnectionProvider;
@@ -90,17 +91,28 @@ public class S3TvfLoad {
                 sqlBuilder.buildInsertSql(
                         database, table, label, objectKeys, columns, 
deleteSignEnabled);
         for (int attempt = 0; attempt <= MAX_INSERT_RETRIES; attempt++) {
+            long insertStartedAtNanos = System.nanoTime();
             try {
                 executeInsert(sql);
-                LOG.info("S3 TVF load committed with label {}", label);
+                LOG.info(
+                        "S3 TVF insert completed, label={}, objectCount={}, 
attempt={}, "
+                                + "insertTimeMs={}",
+                        label,
+                        objectKeys.size(),
+                        attempt + 1,
+                        TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - 
insertStartedAtNanos));
                 return;
             } catch (SQLException e) {
                 LOG.warn(
-                        "S3 TVF insert failed for label {} on attempt {} 
(SQLState={}, errorCode={})",
+                        "S3 TVF insert failed, label={}, objectCount={}, 
attempt={}, "
+                                + "insertTimeMs={}, SQLState={}, errorCode={}",
                         label,
+                        objectKeys.size(),
                         attempt + 1,
+                        TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - 
insertStartedAtNanos),
                         e.getSQLState(),
-                        e.getErrorCode());
+                        e.getErrorCode(),
+                        e);
                 if (isLabelAlreadyUsed(e, label)) {
                     try {
                         if (handleLabelAlreadyUsed(label)) {
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriterTest.java
 
b/src/test/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriterTest.java
index f074b44..fed5c4b 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriterTest.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/writer/AsyncS3TvfWriterTest.java
@@ -29,6 +29,7 @@ import static org.mockito.Mockito.when;
 
 import java.io.IOException;
 import java.io.InputStream;
+import java.io.StringWriter;
 import java.nio.charset.StandardCharsets;
 import java.util.Arrays;
 import java.util.LinkedHashMap;
@@ -53,6 +54,8 @@ import 
org.apache.doris.kafka.connector.service.DorisSystemService;
 import org.apache.doris.kafka.connector.writer.s3.S3ObjectStore;
 import org.apache.doris.kafka.connector.writer.s3.S3TvfLoad;
 import org.apache.kafka.connect.sink.SinkRecord;
+import org.apache.log4j.PatternLayout;
+import org.apache.log4j.WriterAppender;
 import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.ArgumentCaptor;
@@ -60,6 +63,35 @@ import org.mockito.ArgumentCaptor;
 public class AsyncS3TvfWriterTest {
     private static final String LABEL_PREFIX = "tvf_demo_orders_";
 
+    @Test
+    public void testLogsUploadedObjectMetrics() throws Exception {
+        RecordingObjectStore store = new RecordingObjectStore();
+        S3TvfLoad load = mock(S3TvfLoad.class);
+        RecordService records = mock(RecordService.class);
+        SinkRecord record = TestRecordBuffer.newSinkRecord("ignored", 1);
+        
when(records.getProcessedRecord(record)).thenReturn("{\"id\":1,\"name\":\"first\"}");
+        AsyncS3TvfWriter writer = writer(options(1024, 100), store, load, 
records);
+        StringWriter logs = new StringWriter();
+        WriterAppender appender = new WriterAppender(new 
PatternLayout("%m%n"), logs);
+        org.apache.log4j.Logger logger = 
org.apache.log4j.Logger.getLogger(AsyncS3TvfWriter.class);
+        logger.addAppender(appender);
+
+        try {
+            writer.insert(record);
+            writer.commitFlush();
+        } finally {
+            logger.removeAppender(appender);
+            appender.close();
+            writer.close();
+        }
+
+        String output = logs.toString();
+        Assert.assertTrue(output.contains("S3 TVF object upload completed"));
+        
Assert.assertTrue(output.contains("objectKey=objects/tvf/demo_orders/"));
+        Assert.assertTrue(output.contains("sizeBytes=24"));
+        Assert.assertTrue(output.matches("(?s).*uploadTimeMs=\\d+.*"));
+    }
+
     @Test
     public void testUsesTaskIdAndOneBatchLabelForAllFiles() throws Exception {
         DorisOptions options = options(1024, 1);
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStoreTest.java
 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStoreTest.java
index 8d74838..5dbb6c0 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStoreTest.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3ClientObjectStoreTest.java
@@ -23,8 +23,8 @@ import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
-import org.apache.commons.io.IOUtils;
 import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.ArgumentCaptor;
@@ -37,13 +37,14 @@ import 
software.amazon.awssdk.services.s3.model.PutObjectResponse;
 public class S3ClientObjectStoreTest {
 
     @Test
-    public void testPutUsesExactBucketKeyAndJsonLinesContentType() throws 
Exception {
+    public void testPutUsesRepeatableContentProviderWithoutCopying() throws 
Exception {
         S3Client client = Mockito.mock(S3Client.class);
         when(client.putObject(any(PutObjectRequest.class), 
any(RequestBody.class)))
                 .thenReturn(PutObjectResponse.builder().build());
         S3ClientObjectStore store = new S3ClientObjectStore(client, "staging");
+        byte[] content = "{\"id\":1}\n".getBytes(StandardCharsets.UTF_8);
 
-        store.put("kafka/orders/file.json", 
"{\"id\":1}\n".getBytes(StandardCharsets.UTF_8));
+        store.put("kafka/orders/file.json", content);
 
         ArgumentCaptor<PutObjectRequest> request = 
ArgumentCaptor.forClass(PutObjectRequest.class);
         ArgumentCaptor<RequestBody> body = 
ArgumentCaptor.forClass(RequestBody.class);
@@ -51,10 +52,16 @@ public class S3ClientObjectStoreTest {
         Assert.assertEquals("staging", request.getValue().bucket());
         Assert.assertEquals("kafka/orders/file.json", 
request.getValue().key());
         Assert.assertEquals("application/x-ndjson", 
request.getValue().contentType());
-        Assert.assertEquals(
-                "{\"id\":1}\n",
-                IOUtils.toString(
-                        body.getValue().contentStreamProvider().newStream(),
-                        StandardCharsets.UTF_8));
+
+        content[0] = '[';
+        try (InputStream input = 
body.getValue().contentStreamProvider().newStream();
+                InputStream retryInput = 
body.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);
+        }
     }
 }
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoadTest.java 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoadTest.java
index b8ef4ef..49bee30 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoadTest.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfLoadTest.java
@@ -27,6 +27,7 @@ import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.io.StringWriter;
 import java.sql.Connection;
 import java.sql.ResultSet;
 import java.sql.SQLException;
@@ -38,6 +39,9 @@ import java.util.Map;
 import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
 import org.apache.doris.kafka.connector.connection.ConnectionProvider;
 import org.apache.doris.kafka.connector.exception.DorisException;
+import org.apache.log4j.PatternLayout;
+import org.apache.log4j.WriterAppender;
+import org.junit.After;
 import org.junit.Assert;
 import org.junit.Before;
 import org.junit.Test;
@@ -48,6 +52,8 @@ public class S3TvfLoadTest {
     private Connection connection;
     private Statement statement;
     private String insertSql;
+    private StringWriter logs;
+    private WriterAppender appender;
 
     @Before
     public void setUp() throws Exception {
@@ -58,6 +64,39 @@ public class S3TvfLoadTest {
         when(connection.createStatement()).thenReturn(statement);
         insertSql =
                 sqlBuilder().buildInsertSql("demo", "orders", "label", 
files(), columns(), false);
+        logs = new StringWriter();
+        appender = new WriterAppender(new PatternLayout("%m%n"), logs);
+        
org.apache.log4j.Logger.getLogger(S3TvfLoad.class).addAppender(appender);
+    }
+
+    @After
+    public void tearDown() {
+        
org.apache.log4j.Logger.getLogger(S3TvfLoad.class).removeAppender(appender);
+        appender.close();
+    }
+
+    @Test
+    public void testLogsInsertMetrics() {
+        load(Collections.emptyMap()).load("label", files());
+
+        String output = logs.toString();
+        Assert.assertTrue(
+                output.contains("S3 TVF insert completed, label=label, 
objectCount=1, attempt=1"));
+        Assert.assertTrue(output.matches("(?s).*insertTimeMs=\\d+.*"));
+    }
+
+    @Test
+    public void testLogsInsertFailureCause() throws Exception {
+        when(statement.execute(insertSql)).thenThrow(new 
SQLException("temporary"));
+
+        try {
+            load(Collections.emptyMap()).load("label", files());
+            Assert.fail("Expected load failure");
+        } catch (DorisException expected) {
+            // The warning must keep the SQLException stack trace for 
diagnosis.
+        }
+
+        Assert.assertTrue(logs.toString().contains("java.sql.SQLException: 
temporary"));
     }
 
     @Test


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

Reply via email to