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 40b80e9  [feat] Support IAM roles and gzip for S3 TVF sink (#107)
40b80e9 is described below

commit 40b80e9897508baf56728b3757725947d103957c
Author: wudi <[email protected]>
AuthorDate: Fri Sep 18 10:03:51 2026 +0800

    [feat] Support IAM roles and gzip for S3 TVF sink (#107)
---
 pom.xml                                            |  15 ++
 .../doris/kafka/connector/cfg/DorisOptions.java    |   6 +
 .../connector/cfg/DorisSinkConnectorConfig.java    |  24 ++-
 .../doris/kafka/connector/cfg/S3TvfOptions.java    |  55 ++++-
 .../kafka/connector/utils/ConfigCheckUtils.java    |  12 ++
 .../kafka/connector/writer/AsyncS3TvfWriter.java   |  24 ++-
 .../connector/writer/s3/S3ClientObjectStore.java   |  74 ++++++-
 .../doris/kafka/connector/writer/s3/S3TvfLoad.java |   2 +-
 .../kafka/connector/writer/s3/S3TvfSqlBuilder.java |  27 ++-
 .../kafka/connector/cfg/S3TvfOptionsTest.java      |  25 +++
 .../kafka/connector/cfg/TestDorisOptions.java      |  13 ++
 .../cfg/TestDorisSinkConnectorConfig.java          |  18 ++
 .../connector/e2e/sink/S3TvfIamRoleITCase.java     | 226 +++++++++++++++++++++
 .../kafka/connector/e2e/sink/S3TvfSinkITCase.java  |   2 +-
 .../connector/writer/AsyncS3TvfWriterTest.java     |  44 ++++
 .../writer/s3/S3ClientObjectStoreTest.java         |  21 ++
 .../connector/writer/s3/S3TvfSqlBuilderTest.java   |  29 +++
 src/test/resources/s3-tvf-iam-role-sink.properties |  29 +++
 18 files changed, 628 insertions(+), 18 deletions(-)

diff --git a/pom.xml b/pom.xml
index 3b71681..08ae6fa 100644
--- a/pom.xml
+++ b/pom.xml
@@ -174,6 +174,21 @@
                 </exclusion>
             </exclusions>
         </dependency>
+        <dependency>
+            <groupId>software.amazon.awssdk</groupId>
+            <artifactId>sts</artifactId>
+            <version>${awssdk.version}</version>
+            <exclusions>
+                <exclusion>
+                    <groupId>software.amazon.awssdk</groupId>
+                    <artifactId>apache-client</artifactId>
+                </exclusion>
+                <exclusion>
+                    <groupId>software.amazon.awssdk</groupId>
+                    <artifactId>netty-nio-client</artifactId>
+                </exclusion>
+            </exclusions>
+        </dependency>
         <dependency>
             <groupId>software.amazon.awssdk</groupId>
             <artifactId>url-connection-client</artifactId>
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/cfg/DorisOptions.java 
b/src/main/java/org/apache/doris/kafka/connector/cfg/DorisOptions.java
index da95e8d..813814d 100644
--- a/src/main/java/org/apache/doris/kafka/connector/cfg/DorisOptions.java
+++ b/src/main/java/org/apache/doris/kafka/connector/cfg/DorisOptions.java
@@ -232,6 +232,8 @@ public class DorisOptions {
                 .setPrefix(config.get(DorisSinkConnectorConfig.SINK_S3_PREFIX))
                 
.setAccessKey(config.get(DorisSinkConnectorConfig.SINK_S3_ACCESS_KEY))
                 
.setSecretKey(config.get(DorisSinkConnectorConfig.SINK_S3_SECRET_KEY))
+                
.setRoleArn(config.get(DorisSinkConnectorConfig.SINK_S3_ROLE_ARN))
+                
.setExternalId(config.get(DorisSinkConnectorConfig.SINK_S3_EXTERNAL_ID))
                 .setPathStyleAccess(
                         Boolean.parseBoolean(
                                 config.getOrDefault(
@@ -411,6 +413,10 @@ public class DorisOptions {
         return streamLoadProp;
     }
 
+    public boolean isGzipCompressionEnabled() {
+        return 
"gz".equalsIgnoreCase(streamLoadProp.getProperty("compress_type", "gz").trim());
+    }
+
     public String getLabelPrefix() {
         return this.labelPrefix;
     }
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/cfg/DorisSinkConnectorConfig.java
 
b/src/main/java/org/apache/doris/kafka/connector/cfg/DorisSinkConnectorConfig.java
index 19f86e9..ee23d0f 100644
--- 
a/src/main/java/org/apache/doris/kafka/connector/cfg/DorisSinkConnectorConfig.java
+++ 
b/src/main/java/org/apache/doris/kafka/connector/cfg/DorisSinkConnectorConfig.java
@@ -100,6 +100,8 @@ public class DorisSinkConnectorConfig {
     public static final String SINK_S3_PREFIX = "sink.s3.prefix";
     public static final String SINK_S3_ACCESS_KEY = "sink.s3.access-key";
     public static final String SINK_S3_SECRET_KEY = "sink.s3.secret-key";
+    public static final String SINK_S3_ROLE_ARN = "sink.s3.role-arn";
+    public static final String SINK_S3_EXTERNAL_ID = "sink.s3.external-id";
     public static final String SINK_S3_PATH_STYLE_ACCESS = 
"sink.s3.path-style-access";
     public static final boolean SINK_S3_PATH_STYLE_ACCESS_DEFAULT = false;
     public static final String CONVERTER_MODE = "converter.mode";
@@ -485,6 +487,26 @@ public class DorisSinkConnectorConfig {
                         6,
                         ConfigDef.Width.NONE,
                         SINK_S3_SECRET_KEY)
+                .define(
+                        SINK_S3_ROLE_ARN,
+                        Type.STRING,
+                        null,
+                        Importance.HIGH,
+                        "AWS IAM role ARN used to access S3",
+                        TVF_CONFIG,
+                        7,
+                        ConfigDef.Width.NONE,
+                        SINK_S3_ROLE_ARN)
+                .define(
+                        SINK_S3_EXTERNAL_ID,
+                        Type.STRING,
+                        null,
+                        Importance.MEDIUM,
+                        "External ID used when assuming the AWS IAM role",
+                        TVF_CONFIG,
+                        8,
+                        ConfigDef.Width.NONE,
+                        SINK_S3_EXTERNAL_ID)
                 .define(
                         SINK_S3_PATH_STYLE_ACCESS,
                         Type.BOOLEAN,
@@ -492,7 +514,7 @@ public class DorisSinkConnectorConfig {
                         Importance.MEDIUM,
                         "Whether to use S3 path-style access",
                         TVF_CONFIG,
-                        7,
+                        9,
                         ConfigDef.Width.NONE,
                         SINK_S3_PATH_STYLE_ACCESS);
     }
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/cfg/S3TvfOptions.java 
b/src/main/java/org/apache/doris/kafka/connector/cfg/S3TvfOptions.java
index d957d7d..83e9a07 100644
--- a/src/main/java/org/apache/doris/kafka/connector/cfg/S3TvfOptions.java
+++ b/src/main/java/org/apache/doris/kafka/connector/cfg/S3TvfOptions.java
@@ -30,6 +30,8 @@ public class S3TvfOptions {
     private final String prefix;
     private final String accessKey;
     private final String secretKey;
+    private final String roleArn;
+    private final String externalId;
     private final boolean pathStyleAccess;
 
     private S3TvfOptions(Builder builder) {
@@ -37,11 +39,14 @@ public class S3TvfOptions {
         this.region = requireNonEmpty(builder.region, "sink.s3.region");
         this.bucket = requireNonEmpty(builder.bucket, "sink.s3.bucket");
         this.prefix = requireNonEmpty(builder.prefix, "sink.s3.prefix");
-        this.accessKey = requireNonEmpty(builder.accessKey, 
"sink.s3.access-key");
-        this.secretKey = requireNonEmpty(builder.secretKey, 
"sink.s3.secret-key");
+        this.accessKey = trimToNull(builder.accessKey);
+        this.secretKey = trimToNull(builder.secretKey);
+        this.roleArn = trimToNull(builder.roleArn);
+        this.externalId = trimToNull(builder.externalId);
         this.pathStyleAccess = builder.pathStyleAccess;
         validateEndpoint(endpoint);
         validatePrefix(prefix);
+        validateCredentials();
     }
 
     public static Builder builder() {
@@ -72,6 +77,22 @@ public class S3TvfOptions {
         return secretKey;
     }
 
+    public String getRoleArn() {
+        return roleArn;
+    }
+
+    public String getExternalId() {
+        return externalId;
+    }
+
+    public boolean hasRoleArn() {
+        return roleArn != null;
+    }
+
+    public boolean hasStaticCredentials() {
+        return accessKey != null;
+    }
+
     public boolean isPathStyleAccess() {
         return pathStyleAccess;
     }
@@ -103,6 +124,24 @@ public class S3TvfOptions {
         return value.trim();
     }
 
+    private static String trimToNull(String value) {
+        return value == null || value.trim().isEmpty() ? null : value.trim();
+    }
+
+    private void validateCredentials() {
+        if ((accessKey == null) != (secretKey == null)) {
+            throw new IllegalArgumentException(
+                    "sink.s3.access-key and sink.s3.secret-key must be 
configured together");
+        }
+        if (accessKey == null && roleArn == null) {
+            throw new IllegalArgumentException(
+                    "S3 TVF requires either access/secret keys or 
sink.s3.role-arn");
+        }
+        if (externalId != null && roleArn == null) {
+            throw new IllegalArgumentException("sink.s3.external-id requires 
sink.s3.role-arn");
+        }
+    }
+
     private static void validatePrefix(String prefix) {
         for (char character : "*?[]{},\\".toCharArray()) {
             if (prefix.indexOf(character) >= 0) {
@@ -133,6 +172,8 @@ public class S3TvfOptions {
         private String prefix;
         private String accessKey;
         private String secretKey;
+        private String roleArn;
+        private String externalId;
         private boolean pathStyleAccess;
 
         public Builder setEndpoint(String endpoint) {
@@ -165,6 +206,16 @@ public class S3TvfOptions {
             return this;
         }
 
+        public Builder setRoleArn(String roleArn) {
+            this.roleArn = roleArn;
+            return this;
+        }
+
+        public Builder setExternalId(String externalId) {
+            this.externalId = externalId;
+            return this;
+        }
+
         public Builder setPathStyleAccess(boolean pathStyleAccess) {
             this.pathStyleAccess = pathStyleAccess;
             return this;
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/utils/ConfigCheckUtils.java 
b/src/main/java/org/apache/doris/kafka/connector/utils/ConfigCheckUtils.java
index e308990..0775566 100644
--- a/src/main/java/org/apache/doris/kafka/connector/utils/ConfigCheckUtils.java
+++ b/src/main/java/org/apache/doris/kafka/connector/utils/ConfigCheckUtils.java
@@ -251,8 +251,20 @@ public class ConfigCheckUtils {
                             
.setPrefix(config.get(DorisSinkConnectorConfig.SINK_S3_PREFIX))
                             
.setAccessKey(config.get(DorisSinkConnectorConfig.SINK_S3_ACCESS_KEY))
                             
.setSecretKey(config.get(DorisSinkConnectorConfig.SINK_S3_SECRET_KEY))
+                            
.setRoleArn(config.get(DorisSinkConnectorConfig.SINK_S3_ROLE_ARN))
+                            
.setExternalId(config.get(DorisSinkConnectorConfig.SINK_S3_EXTERNAL_ID))
                             
.setPathStyleAccess(Boolean.parseBoolean(pathStyleAccess))
                             .build();
+                    String compressType =
+                            config.getOrDefault(
+                                            
DorisSinkConnectorConfig.STREAM_LOAD_PROP_PREFIX
+                                                    + "compress_type",
+                                            "gz")
+                                    .trim();
+                    if (!compressType.isEmpty() && 
!"gz".equalsIgnoreCase(compressType)) {
+                        throw new IllegalArgumentException(
+                                "TVF write mode only supports 'gz' or an empty 
compress_type");
+                    }
                     TvfColumnUtils.resolveColumns(
                             config.get(
                                     
DorisSinkConnectorConfig.STREAM_LOAD_PROP_PREFIX + "columns"));
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 19c8ca2..3e4bb22 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
@@ -20,6 +20,7 @@
 package org.apache.doris.kafka.connector.writer;
 
 import java.io.ByteArrayOutputStream;
+import java.io.IOException;
 import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.Collections;
@@ -32,6 +33,7 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
+import java.util.zip.GZIPOutputStream;
 import org.apache.doris.kafka.connector.cfg.DorisOptions;
 import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
 import org.apache.doris.kafka.connector.connection.ConnectionProvider;
@@ -224,8 +226,15 @@ public class AsyncS3TvfWriter extends DorisWriter {
         startBatchIfNeeded();
         int currentFileNumber = fileNumber++;
         String label = buildLabel();
+        boolean gzipEnabled = dorisOptions.isGzipCompressionEnabled();
         String fileName =
-                label + "_" + dorisOptions.getTaskId() + "_" + 
currentFileNumber + ".json";
+                label
+                        + "_"
+                        + dorisOptions.getTaskId()
+                        + "_"
+                        + currentFileNumber
+                        + ".json"
+                        + (gzipEnabled ? ".gz" : "");
         String objectKey = buildObjectKey(fileName);
         byte[] content = tvfBuffer.toByteArray();
         int recordCount = bufferedRecords;
@@ -236,14 +245,15 @@ public class AsyncS3TvfWriter extends DorisWriter {
                     }
                     long uploadStartedAtNanos = System.nanoTime();
                     try {
-                        objectStore.put(objectKey, content);
+                        byte[] uploadContent = gzipEnabled ? gzip(content) : 
content;
+                        objectStore.put(objectKey, uploadContent);
                         uploadedObjectKeys.add(objectKey);
                         LOG.info(
                                 "S3 TVF object upload completed, fileName={}, 
objectKey={}, "
                                         + "sizeBytes={}, uploadTimeMs={}",
                                 fileName,
                                 objectKey,
-                                content.length,
+                                uploadContent.length,
                                 TimeUnit.NANOSECONDS.toMillis(
                                         System.nanoTime() - 
uploadStartedAtNanos));
                     } catch (Exception e) {
@@ -274,6 +284,14 @@ public class AsyncS3TvfWriter extends DorisWriter {
                 recordCount);
     }
 
+    private static byte[] gzip(byte[] content) throws IOException {
+        ByteArrayOutputStream output = new ByteArrayOutputStream();
+        try (GZIPOutputStream gzip = new GZIPOutputStream(output)) {
+            gzip.write(content);
+        }
+        return output.toByteArray();
+    }
+
     private void runUploadLoop() {
         LOG.info("S3 TVF upload worker started");
         try {
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 bb4aa80..d35f6ba 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
@@ -24,6 +24,8 @@ import java.io.IOException;
 import java.net.URI;
 import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
 import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
+import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider;
+import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;
 import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
 import software.amazon.awssdk.core.sync.RequestBody;
 import software.amazon.awssdk.http.urlconnection.UrlConnectionHttpClient;
@@ -31,16 +33,24 @@ import software.amazon.awssdk.regions.Region;
 import software.amazon.awssdk.services.s3.S3Client;
 import software.amazon.awssdk.services.s3.S3Configuration;
 import software.amazon.awssdk.services.s3.model.PutObjectRequest;
+import software.amazon.awssdk.services.sts.StsClient;
+import 
software.amazon.awssdk.services.sts.auth.StsAssumeRoleCredentialsProvider;
+import software.amazon.awssdk.services.sts.model.AssumeRoleRequest;
 
 /** AWS SDK based object store used by the S3 TVF writer. */
 public class S3ClientObjectStore implements S3ObjectStore {
     private static final String JSON_LINES_CONTENT_TYPE = 
"application/x-ndjson";
+    private static final String ROLE_SESSION_NAME = "doris-kafka-connector";
 
     private final S3Client s3Client;
     private final String bucket;
+    private DefaultCredentialsProvider defaultCredentialsProvider;
+    private StsClient stsClient;
+    private StsAssumeRoleCredentialsProvider assumeRoleCredentialsProvider;
 
     public S3ClientObjectStore(S3TvfOptions options) {
-        this(createClient(options), options.getBucket());
+        this.bucket = options.getBucket();
+        this.s3Client = createClient(options);
     }
 
     S3ClientObjectStore(S3Client s3Client, String bucket) {
@@ -70,17 +80,26 @@ public class S3ClientObjectStore implements S3ObjectStore {
 
     @Override
     public void close() {
-        s3Client.close();
+        try {
+            s3Client.close();
+        } finally {
+            if (assumeRoleCredentialsProvider != null) {
+                assumeRoleCredentialsProvider.close();
+            }
+            if (stsClient != null) {
+                stsClient.close();
+            }
+            if (defaultCredentialsProvider != null) {
+                defaultCredentialsProvider.close();
+            }
+        }
     }
 
-    private static S3Client createClient(S3TvfOptions options) {
+    private S3Client createClient(S3TvfOptions options) {
         return S3Client.builder()
                 .endpointOverride(URI.create(options.getEndpoint()))
                 .region(Region.of(options.getRegion()))
-                .credentialsProvider(
-                        StaticCredentialsProvider.create(
-                                AwsBasicCredentials.create(
-                                        options.getAccessKey(), 
options.getSecretKey())))
+                .credentialsProvider(createCredentialsProvider(options))
                 .httpClientBuilder(UrlConnectionHttpClient.builder())
                 .serviceConfiguration(
                         S3Configuration.builder()
@@ -88,4 +107,45 @@ public class S3ClientObjectStore implements S3ObjectStore {
                                 .build())
                 .build();
     }
+
+    private AwsCredentialsProvider createCredentialsProvider(S3TvfOptions 
options) {
+        if (!options.hasRoleArn()) {
+            return staticCredentialsProvider(options);
+        }
+        AwsCredentialsProvider sourceCredentialsProvider;
+        if (options.hasStaticCredentials()) {
+            sourceCredentialsProvider = staticCredentialsProvider(options);
+        } else {
+            defaultCredentialsProvider = 
DefaultCredentialsProvider.builder().build();
+            sourceCredentialsProvider = defaultCredentialsProvider;
+        }
+        stsClient =
+                StsClient.builder()
+                        .region(Region.of(options.getRegion()))
+                        .credentialsProvider(sourceCredentialsProvider)
+                        .httpClientBuilder(UrlConnectionHttpClient.builder())
+                        .build();
+        assumeRoleCredentialsProvider =
+                StsAssumeRoleCredentialsProvider.builder()
+                        .stsClient(stsClient)
+                        .refreshRequest(buildAssumeRoleRequest(options))
+                        .build();
+        return assumeRoleCredentialsProvider;
+    }
+
+    static AssumeRoleRequest buildAssumeRoleRequest(S3TvfOptions options) {
+        AssumeRoleRequest.Builder request =
+                AssumeRoleRequest.builder()
+                        .roleArn(options.getRoleArn())
+                        .roleSessionName(ROLE_SESSION_NAME);
+        if (options.getExternalId() != null) {
+            request.externalId(options.getExternalId());
+        }
+        return request.build();
+    }
+
+    private static StaticCredentialsProvider 
staticCredentialsProvider(S3TvfOptions options) {
+        return StaticCredentialsProvider.create(
+                AwsBasicCredentials.create(options.getAccessKey(), 
options.getSecretKey()));
+    }
 }
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 1f3942e..d79a9ce 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
@@ -61,7 +61,7 @@ public class S3TvfLoad {
             String table) {
         this(
                 connectionProvider,
-                new S3TvfSqlBuilder(options.getS3TvfOptions()),
+                new S3TvfSqlBuilder(options.getS3TvfOptions(), 
options.isGzipCompressionEnabled()),
                 database,
                 table,
                 options.getTvfColumns(),
diff --git 
a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilder.java 
b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilder.java
index 5f0fa25..23c7421 100644
--- 
a/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilder.java
+++ 
b/src/main/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilder.java
@@ -31,9 +31,15 @@ import org.apache.doris.kafka.connector.writer.LoadConstants;
 /** Builds explicit INSERT SELECT statements for staged S3 TVF files. */
 public class S3TvfSqlBuilder {
     private final S3TvfOptions options;
+    private final boolean gzipEnabled;
 
     public S3TvfSqlBuilder(S3TvfOptions options) {
+        this(options, false);
+    }
+
+    public S3TvfSqlBuilder(S3TvfOptions options, boolean gzipEnabled) {
         this.options = options;
+        this.gzipEnabled = gzipEnabled;
     }
 
     public String buildInsertSql(
@@ -52,6 +58,7 @@ public class S3TvfSqlBuilder {
         }
         String columnSql = joinIdentifiers(loadColumns);
         String uri = buildUri(objectKeys);
+        String credentials = buildCredentials();
         return "INSERT INTO "
                 + quoteIdentifier(database)
                 + "."
@@ -65,9 +72,7 @@ public class S3TvfSqlBuilder {
                 + " FROM S3("
                 + property("uri", uri)
                 + ","
-                + property("s3.access_key", options.getAccessKey())
-                + ","
-                + property("s3.secret_key", options.getSecretKey())
+                + credentials
                 + ","
                 + property("s3.region", options.getRegion())
                 + ","
@@ -76,11 +81,27 @@ public class S3TvfSqlBuilder {
                 + property("format", "json")
                 + ","
                 + property("read_json_by_line", "true")
+                + (gzipEnabled ? "," + property("compress_type", "gz") : "")
                 + ","
                 + property("use_path_style", 
Boolean.toString(options.isPathStyleAccess()))
                 + ")";
     }
 
+    private String buildCredentials() {
+        StringJoiner credentials = new StringJoiner(",");
+        if (options.hasStaticCredentials()) {
+            credentials.add(property("s3.access_key", options.getAccessKey()));
+            credentials.add(property("s3.secret_key", options.getSecretKey()));
+        }
+        if (options.hasRoleArn()) {
+            credentials.add(property("s3.role_arn", options.getRoleArn()));
+            if (options.getExternalId() != null) {
+                credentials.add(property("s3.external_id", 
options.getExternalId()));
+            }
+        }
+        return credentials.toString();
+    }
+
     private String buildUri(List<String> objectKeys) {
         if (objectKeys.size() == 1) {
             return "s3://" + options.getBucket() + "/" + objectKeys.get(0);
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/cfg/S3TvfOptionsTest.java 
b/src/test/java/org/apache/doris/kafka/connector/cfg/S3TvfOptionsTest.java
index ad3392e..f907f8b 100644
--- a/src/test/java/org/apache/doris/kafka/connector/cfg/S3TvfOptionsTest.java
+++ b/src/test/java/org/apache/doris/kafka/connector/cfg/S3TvfOptionsTest.java
@@ -34,6 +34,8 @@ public class S3TvfOptionsTest {
                         .setPrefix("kafka/orders")
                         .setAccessKey("access-key")
                         .setSecretKey("secret-key")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
                         .setPathStyleAccess(true)
                         .build();
 
@@ -43,9 +45,12 @@ public class S3TvfOptionsTest {
         Assert.assertEquals("kafka/orders", options.getPrefix());
         Assert.assertEquals("access-key", options.getAccessKey());
         Assert.assertEquals("secret-key", options.getSecretKey());
+        Assert.assertEquals("arn:aws:iam::123456789012:role/doris", 
options.getRoleArn());
+        Assert.assertEquals("external-id", options.getExternalId());
         Assert.assertTrue(options.isPathStyleAccess());
         Assert.assertFalse(options.toString().contains("access-key"));
         Assert.assertFalse(options.toString().contains("secret-key"));
+        Assert.assertFalse(options.toString().contains("external-id"));
     }
 
     @Test
@@ -65,6 +70,26 @@ public class S3TvfOptionsTest {
         validBuilder().setEndpoint("s3.example.com").build();
     }
 
+    @Test(expected = IllegalArgumentException.class)
+    public void testRequiresCredentialsOrRole() {
+        S3TvfOptions.builder()
+                .setEndpoint("https://s3.example.com";)
+                .setRegion("us-east-1")
+                .setBucket("staging")
+                .setPrefix("kafka/orders")
+                .build();
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testRequiresCompleteStaticCredentials() {
+        validBuilder().setSecretKey(null).build();
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testExternalIdRequiresRole() {
+        validBuilder().setExternalId("external-id").build();
+    }
+
     private static S3TvfOptions.Builder validBuilder() {
         return S3TvfOptions.builder()
                 .setEndpoint("https://s3.example.com";)
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisOptions.java 
b/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisOptions.java
index 1b72824..b3cf0ca 100644
--- a/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisOptions.java
+++ b/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisOptions.java
@@ -154,5 +154,18 @@ public class TestDorisOptions {
                 "ERROR", 
options.getSessionVariables().get("partial_update_new_key_behavior"));
         
Assert.assertFalse(options.getSessionVariables().containsKey("format"));
         
Assert.assertFalse(options.getSessionVariables().containsKey("compress_type"));
+        Assert.assertTrue(options.isGzipCompressionEnabled());
+
+        config.put("sink.properties.compress_type", "");
+        Assert.assertFalse(new 
DorisOptions(config).isGzipCompressionEnabled());
+
+        config.remove(DorisSinkConnectorConfig.SINK_S3_ACCESS_KEY);
+        config.remove(DorisSinkConnectorConfig.SINK_S3_SECRET_KEY);
+        config.put(
+                DorisSinkConnectorConfig.SINK_S3_ROLE_ARN, 
"arn:aws:iam::123456789012:role/doris");
+        config.put(DorisSinkConnectorConfig.SINK_S3_EXTERNAL_ID, 
"external-id");
+        S3TvfOptions roleOptions = new DorisOptions(config).getS3TvfOptions();
+        Assert.assertEquals("arn:aws:iam::123456789012:role/doris", 
roleOptions.getRoleArn());
+        Assert.assertEquals("external-id", roleOptions.getExternalId());
     }
 }
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisSinkConnectorConfig.java
 
b/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisSinkConnectorConfig.java
index da05d62..b829529 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisSinkConnectorConfig.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/cfg/TestDorisSinkConnectorConfig.java
@@ -267,6 +267,24 @@ public class TestDorisSinkConnectorConfig {
         ConfigCheckUtils.validateConfig(getS3TvfConfig());
     }
 
+    @Test
+    public void testS3TvfIamRoleConfig() {
+        Map<String, String> config = getS3TvfConfig();
+        config.remove(DorisSinkConnectorConfig.SINK_S3_ACCESS_KEY);
+        config.remove(DorisSinkConnectorConfig.SINK_S3_SECRET_KEY);
+        config.put(
+                DorisSinkConnectorConfig.SINK_S3_ROLE_ARN, 
"arn:aws:iam::123456789012:role/doris");
+        config.put(DorisSinkConnectorConfig.SINK_S3_EXTERNAL_ID, 
"external-id");
+        ConfigCheckUtils.validateConfig(config);
+    }
+
+    @Test(expected = DorisException.class)
+    public void testS3TvfRejectsUnsupportedCompression() {
+        Map<String, String> config = getS3TvfConfig();
+        config.put(DorisSinkConnectorConfig.STREAM_LOAD_PROP_PREFIX + 
"compress_type", "zstd");
+        ConfigCheckUtils.validateConfig(config);
+    }
+
     @Test(expected = DorisException.class)
     public void testRejectLegacyS3TvfLoadModel() {
         Map<String, String> config = getS3TvfConfig();
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfIamRoleITCase.java
 
b/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfIamRoleITCase.java
new file mode 100644
index 0000000..3cef9e8
--- /dev/null
+++ 
b/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfIamRoleITCase.java
@@ -0,0 +1,226 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.doris.kafka.connector.e2e.sink;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.node.ObjectNode;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.HttpURLConnection;
+import java.net.URL;
+import java.nio.charset.StandardCharsets;
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.Properties;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+import org.apache.doris.kafka.connector.e2e.doris.DorisCustomerServiceImpl;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.junit.Assert;
+import org.junit.Assume;
+import org.junit.BeforeClass;
+import org.junit.Test;
+
+/** Opt-in integration test for Kafka Connect S3 TVF writes with an AWS IAM 
role. */
+public class S3TvfIamRoleITCase {
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+    private static final String DATABASE = "test_s3_tvf_iam_role";
+    private static DorisCustomerServiceImpl doris;
+
+    @BeforeClass
+    public static void useExternalEnvironment() {
+        Assume.assumeTrue(Boolean.getBoolean("s3_tvf_iam_role_it"));
+        Assume.assumeTrue(Boolean.getBoolean("customer_env"));
+        requiredProperty("kafka_bootstrap_servers");
+        requiredProperty("kafka_connect_url");
+        requiredProperty("s3_endpoint");
+        requiredProperty("s3_region");
+        requiredProperty("s3_bucket");
+        requiredProperty("s3_role_arn");
+        doris = new DorisCustomerServiceImpl();
+        doris.startContainer();
+    }
+
+    @Test
+    public void testWritesThroughIamRole() throws Exception {
+        String suffix = UUID.randomUUID().toString().replace("-", "");
+        String connector = "s3-tvf-iam-role-" + suffix;
+        String topic = "s3-tvf-iam-role-" + suffix;
+        String table = "iam_role_" + suffix;
+        boolean registered = false;
+        try {
+            executeSql(
+                    "CREATE DATABASE IF NOT EXISTS `" + DATABASE + "`",
+                    "CREATE TABLE `"
+                            + DATABASE
+                            + "`.`"
+                            + table
+                            + "` (`id` INT, `name` VARCHAR(64)) "
+                            + "DUPLICATE KEY(`id`) DISTRIBUTED BY HASH(`id`) 
BUCKETS 1 "
+                            + "PROPERTIES (\"replication_num\" = \"1\")");
+            registerConnector(connector, topic, table, suffix);
+            registered = true;
+            produce(topic, "{\"id\":1,\"name\":\"kafka\"}");
+            waitForRow(table);
+        } finally {
+            if (registered) {
+                deleteConnector(connector);
+            }
+            executeSql("DROP TABLE IF EXISTS `" + DATABASE + "`.`" + table + 
"`");
+        }
+    }
+
+    private static void registerConnector(
+            String connector, String topic, String table, String suffix) 
throws Exception {
+        Properties properties = new Properties();
+        try (InputStream stream =
+                S3TvfIamRoleITCase.class
+                        .getClassLoader()
+                        
.getResourceAsStream("s3-tvf-iam-role-sink.properties")) {
+            properties.load(stream);
+        }
+        properties.put("topics", topic);
+        properties.put("doris.topic2table.map", topic + ":" + table);
+        properties.put("label.prefix", "iam_role_" + suffix);
+        properties.put("doris.urls", requiredProperty("doris_host"));
+        properties.put("doris.http.port", requiredProperty("doris_http_port"));
+        properties.put("doris.query.port", 
requiredProperty("doris_query_port"));
+        properties.put("doris.user", requiredProperty("doris_user"));
+        properties.put("doris.password", System.getProperty("doris_passwd", 
""));
+        properties.put("doris.database", DATABASE);
+        properties.put("sink.s3.endpoint", requiredProperty("s3_endpoint"));
+        properties.put("sink.s3.region", requiredProperty("s3_region"));
+        properties.put("sink.s3.bucket", requiredProperty("s3_bucket"));
+        properties.put(
+                "sink.s3.prefix", System.getProperty("s3_prefix", 
"doris-kafka-connector-it"));
+        properties.put("sink.s3.role-arn", requiredProperty("s3_role_arn"));
+        optionalProperty("s3_external_id")
+                .ifPresent(value -> properties.put("sink.s3.external-id", 
value));
+
+        ObjectNode root = OBJECT_MAPPER.createObjectNode();
+        root.put("name", connector);
+        ObjectNode config = root.putObject("config");
+        properties.forEach((key, value) -> config.put(key.toString(), 
value.toString()));
+        request("POST", "/connectors", OBJECT_MAPPER.writeValueAsBytes(root), 
201);
+    }
+
+    private static void produce(String topic, String value) throws Exception {
+        Properties properties = new Properties();
+        properties.put(
+                ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
+                requiredProperty("kafka_bootstrap_servers"));
+        properties.put(
+                ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        properties.put(
+                ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        try (KafkaProducer<String, String> producer = new 
KafkaProducer<>(properties)) {
+            producer.send(new ProducerRecord<>(topic, value)).get(30, 
TimeUnit.SECONDS);
+        }
+    }
+
+    private static void waitForRow(String table) throws Exception {
+        long deadline = System.nanoTime() + TimeUnit.MINUTES.toNanos(2);
+        while (System.nanoTime() < deadline) {
+            try (Connection connection = doris.getQueryConnection();
+                    Statement statement = connection.createStatement();
+                    ResultSet result =
+                            statement.executeQuery(
+                                    "SELECT id,name FROM `" + DATABASE + "`.`" 
+ table + "`")) {
+                if (result.next()) {
+                    Assert.assertEquals(1, result.getInt(1));
+                    Assert.assertEquals("kafka", result.getString(2));
+                    return;
+                }
+            }
+            Thread.sleep(1000);
+        }
+        Assert.fail("Timed out waiting for Kafka Connect IAM role row");
+    }
+
+    private static void executeSql(String... sql) throws Exception {
+        try (Connection connection = doris.getQueryConnection();
+                Statement statement = connection.createStatement()) {
+            for (String value : sql) {
+                statement.execute(value);
+            }
+        }
+    }
+
+    private static void deleteConnector(String connector) throws Exception {
+        request("DELETE", "/connectors/" + connector, null, 204);
+    }
+
+    private static void request(String method, String path, byte[] body, int 
expected)
+            throws Exception {
+        HttpURLConnection connection =
+                (HttpURLConnection)
+                        new URL(requiredProperty("kafka_connect_url") + 
path).openConnection();
+        connection.setRequestMethod(method);
+        connection.setConnectTimeout(10000);
+        connection.setReadTimeout(30000);
+        if (body != null) {
+            connection.setDoOutput(true);
+            connection.setRequestProperty("Content-Type", "application/json");
+            try (OutputStream output = connection.getOutputStream()) {
+                output.write(body);
+            }
+        }
+        int status = connection.getResponseCode();
+        if (status != expected) {
+            InputStream response =
+                    status >= 400 ? connection.getErrorStream() : 
connection.getInputStream();
+            String message =
+                    response == null ? "" : new String(readAll(response), 
StandardCharsets.UTF_8);
+            throw new IllegalStateException("Kafka Connect returned " + status 
+ ": " + message);
+        }
+        connection.disconnect();
+    }
+
+    private static byte[] readAll(InputStream input) throws Exception {
+        byte[] buffer = new byte[1024];
+        java.io.ByteArrayOutputStream output = new 
java.io.ByteArrayOutputStream();
+        try (InputStream stream = input) {
+            int length;
+            while ((length = stream.read(buffer)) != -1) {
+                output.write(buffer, 0, length);
+            }
+        }
+        return output.toByteArray();
+    }
+
+    private static String requiredProperty(String name) {
+        return optionalProperty(name)
+                .orElseThrow(
+                        () ->
+                                new IllegalArgumentException(
+                                        "Missing required system property: " + 
name));
+    }
+
+    private static java.util.Optional<String> optionalProperty(String name) {
+        String value = System.getProperty(name);
+        return value == null || value.trim().isEmpty()
+                ? java.util.Optional.empty()
+                : java.util.Optional.of(value.trim());
+    }
+}
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfSinkITCase.java 
b/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfSinkITCase.java
index 18f5a04..73d0f35 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfSinkITCase.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/e2e/sink/S3TvfSinkITCase.java
@@ -455,7 +455,7 @@ public class S3TvfSinkITCase extends 
AbstractStringE2ESinkTest {
                         "^"
                                 + Pattern.quote(
                                         JSON_LABEL_PREFIX + "_" + DATABASE + 
"_" + JSON_TABLE + "_")
-                                + "([0-9a-f]{32})_0_([0-9]+)\\.json$");
+                                + "([0-9a-f]{32})_0_([0-9]+)\\.json\\.gz$");
         Set<String> batchUuids = new HashSet<>();
         Set<String> fileNumbers = new HashSet<>();
         for (String objectKey : objectKeys) {
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 fed5c4b..d5070de 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
@@ -27,6 +27,8 @@ import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.io.InputStream;
 import java.io.StringWriter;
@@ -44,6 +46,7 @@ import java.util.concurrent.Future;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.zip.GZIPInputStream;
 import org.apache.doris.kafka.connector.cfg.DorisOptions;
 import org.apache.doris.kafka.connector.cfg.DorisSinkConnectorConfig;
 import org.apache.doris.kafka.connector.connection.ConnectionProvider;
@@ -125,6 +128,26 @@ public class AsyncS3TvfWriterTest {
         writer.close();
     }
 
+    @Test
+    public void testGzipUpload() 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, "tvf", true), 
store, load, records);
+
+        writer.insert(record);
+        writer.commitFlush();
+
+        Map.Entry<String, byte[]> object = 
store.objects.entrySet().iterator().next();
+        Assert.assertTrue(object.getKey().endsWith(".json.gz"));
+        Assert.assertEquals(
+                "{\"id\":1,\"name\":\"first\"}\n",
+                new String(gunzip(object.getValue()), StandardCharsets.UTF_8));
+        writer.close();
+    }
+
     @Test
     public void testLabelDoesNotExceedDorisLimit() throws Exception {
         RecordingObjectStore store = new RecordingObjectStore();
@@ -441,6 +464,12 @@ public class AsyncS3TvfWriterTest {
 
     private static DorisOptions options(int bufferSize, int recordCount, 
String labelPrefix)
             throws IOException {
+        return options(bufferSize, recordCount, labelPrefix, false);
+    }
+
+    private static DorisOptions options(
+            int bufferSize, int recordCount, String labelPrefix, boolean 
gzipEnabled)
+            throws IOException {
         InputStream stream =
                 AsyncS3TvfWriterTest.class
                         .getClassLoader()
@@ -464,9 +493,24 @@ public class AsyncS3TvfWriterTest {
         properties.put(DorisSinkConnectorConfig.SINK_S3_ACCESS_KEY, 
"access-key");
         properties.put(DorisSinkConnectorConfig.SINK_S3_SECRET_KEY, 
"secret-key");
         properties.put(DorisSinkConnectorConfig.STREAM_LOAD_PROP_PREFIX + 
"columns", "id,name");
+        properties.put(
+                DorisSinkConnectorConfig.STREAM_LOAD_PROP_PREFIX + 
"compress_type",
+                gzipEnabled ? "gz" : "");
         return new DorisOptions((Map) properties);
     }
 
+    private static byte[] gunzip(byte[] content) throws IOException {
+        ByteArrayOutputStream output = new ByteArrayOutputStream();
+        byte[] buffer = new byte[1024];
+        try (GZIPInputStream input = new GZIPInputStream(new 
ByteArrayInputStream(content))) {
+            int length;
+            while ((length = input.read(buffer)) != -1) {
+                output.write(buffer, 0, length);
+            }
+        }
+        return output.toByteArray();
+    }
+
     private static class RecordingObjectStore implements S3ObjectStore {
         protected final Map<String, byte[]> objects = new LinkedHashMap<>();
         private IOException putFailure;
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 5dbb6c0..eea5bc2 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
@@ -25,6 +25,7 @@ import static org.mockito.Mockito.when;
 
 import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
+import org.apache.doris.kafka.connector.cfg.S3TvfOptions;
 import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.ArgumentCaptor;
@@ -33,9 +34,29 @@ import software.amazon.awssdk.core.sync.RequestBody;
 import software.amazon.awssdk.services.s3.S3Client;
 import software.amazon.awssdk.services.s3.model.PutObjectRequest;
 import software.amazon.awssdk.services.s3.model.PutObjectResponse;
+import software.amazon.awssdk.services.sts.model.AssumeRoleRequest;
 
 public class S3ClientObjectStoreTest {
 
+    @Test
+    public void testBuildAssumeRoleRequest() {
+        S3TvfOptions options =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.example.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("staging")
+                        .setPrefix("kafka/orders")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
+                        .build();
+
+        AssumeRoleRequest request = 
S3ClientObjectStore.buildAssumeRoleRequest(options);
+
+        Assert.assertEquals("arn:aws:iam::123456789012:role/doris", 
request.roleArn());
+        Assert.assertEquals("external-id", request.externalId());
+        Assert.assertEquals("doris-kafka-connector", 
request.roleSessionName());
+    }
+
     @Test
     public void testPutUsesRepeatableContentProviderWithoutCopying() throws 
Exception {
         S3Client client = Mockito.mock(S3Client.class);
diff --git 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilderTest.java
 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilderTest.java
index a3954ae..da3aa1d 100644
--- 
a/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilderTest.java
+++ 
b/src/test/java/org/apache/doris/kafka/connector/writer/s3/S3TvfSqlBuilderTest.java
@@ -72,6 +72,35 @@ public class S3TvfSqlBuilderTest {
         Assert.assertFalse(builder.toString().contains("s\\k"));
     }
 
+    @Test
+    public void testBuildInsertWithIamRoleAndGzip() {
+        S3TvfOptions options =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.example.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("staging")
+                        .setPrefix("objects")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
+                        .build();
+
+        String sql =
+                new S3TvfSqlBuilder(options, true)
+                        .buildInsertSql(
+                                "demo",
+                                "orders",
+                                "label",
+                                Arrays.asList("objects/file.json.gz"),
+                                Arrays.asList("id"),
+                                false);
+
+        Assert.assertTrue(sql.contains("'s3.role_arn' = 
'arn:aws:iam::123456789012:role/doris'"));
+        Assert.assertTrue(sql.contains("'s3.external_id' = 'external-id'"));
+        Assert.assertTrue(sql.contains("'compress_type' = 'gz'"));
+        Assert.assertFalse(sql.contains("s3.access_key"));
+        Assert.assertFalse(sql.contains("s3.secret_key"));
+    }
+
     private static S3TvfOptions options(String accessKey, String secretKey) {
         return S3TvfOptions.builder()
                 .setEndpoint("https://s3.example.com";)
diff --git a/src/test/resources/s3-tvf-iam-role-sink.properties 
b/src/test/resources/s3-tvf-iam-role-sink.properties
new file mode 100644
index 0000000..a20791f
--- /dev/null
+++ b/src/test/resources/s3-tvf-iam-role-sink.properties
@@ -0,0 +1,29 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
+tasks.max=1
+buffer.count.records=1
+buffer.flush.time=1
+buffer.size.bytes=1
+load.model=tvf
+enable.combine.flush=true
+delivery.guarantee=at_least_once
+converter.mode=normal
+key.converter=org.apache.kafka.connect.storage.StringConverter
+value.converter=org.apache.kafka.connect.storage.StringConverter
+sink.properties.columns=id,name


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

Reply via email to