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 523ed9d7 [Feature] Support IAM roles and gzip compression for S3 TVF 
sink (#698)
523ed9d7 is described below

commit 523ed9d7eda962861d9a8d0a97534df6e337183c
Author: wudi <[email protected]>
AuthorDate: Thu Sep 17 09:55:18 2026 +0800

    [Feature] Support IAM roles and gzip compression for S3 TVF sink (#698)
    
    1. Support AWS IAM role authentication for S3 TVF writes, with optional 
external ID and static or default source credentials.
    2. Enable gzip compression by default for TVF staging files, aligned with 
Stream Load compress_type configuration; an empty value disables compression.
    3. Add unit coverage, MinIO compression coverage, and an opt-in real AWS 
IAM role integration case.
---
 .../flink-doris-connector-base/pom.xml             |  15 ++
 .../doris/flink/cfg/DorisExecutionOptions.java     |  17 +-
 .../org/apache/doris/flink/cfg/S3TvfOptions.java   |  65 ++++++-
 .../flink/sink/writer/tvf/S3ClientObjectStore.java |  86 ++++++++-
 .../flink/sink/writer/tvf/S3TvfCommitter.java      |  13 +-
 .../flink/sink/writer/tvf/S3TvfSqlBuilder.java     |  25 ++-
 .../doris/flink/sink/writer/tvf/S3TvfWriter.java   |  28 ++-
 .../doris/flink/table/DorisConfigOptions.java      |  55 +++++-
 .../doris/flink/cfg/DorisExecutionOptionsTest.java |  36 +++-
 .../apache/doris/flink/cfg/S3TvfOptionsTest.java   |  32 +++-
 .../sink/writer/tvf/S3ClientObjectStoreTest.java   |  17 ++
 .../flink/sink/writer/tvf/S3TvfCommitterTest.java  |  18 +-
 .../flink/sink/writer/tvf/S3TvfSqlBuilderTest.java |  53 +++++-
 .../flink/sink/writer/tvf/S3TvfWriterTest.java     |  42 ++++-
 .../flink/sink/writer/tvf/S3TvfWriterAdapter.java  |   3 +-
 .../flink/table/DorisDynamicTableFactory.java      |   4 +
 .../flink/table/DorisDynamicTableFactoryTest.java  |  46 ++++-
 .../flink/sink/writer/tvf/S3TvfWriterAdapter.java  |   3 +-
 .../flink/table/DorisDynamicTableFactory.java      |   4 +
 .../flink/table/DorisDynamicTableFactoryTest.java  |  46 ++++-
 .../doris/flink/sink/S3TvfIamRoleITCase.java       | 203 +++++++++++++++++++++
 .../apache/doris/flink/sink/S3TvfSinkITCase.java   |  81 +++++++-
 22 files changed, 837 insertions(+), 55 deletions(-)

diff --git a/flink-doris-connector/flink-doris-connector-base/pom.xml 
b/flink-doris-connector/flink-doris-connector-base/pom.xml
index 8f1da7e7..c2439153 100644
--- a/flink-doris-connector/flink-doris-connector-base/pom.xml
+++ b/flink-doris-connector/flink-doris-connector-base/pom.xml
@@ -172,6 +172,21 @@ under the License.
                 </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/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 138e7cae..f0952e40 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
@@ -221,6 +221,12 @@ public class DorisExecutionOptions implements Serializable 
{
         return streamLoadProp;
     }
 
+    public boolean isGzipCompressionEnabled() {
+        return streamLoadProp == null
+                || COMPRESS_TYPE_GZ.equalsIgnoreCase(
+                        streamLoadProp.getProperty(COMPRESS_TYPE, 
COMPRESS_TYPE_GZ).trim());
+    }
+
     public Boolean getDeletable() {
         return enableDelete;
     }
@@ -601,12 +607,17 @@ public class DorisExecutionOptions implements 
Serializable {
             }
 
             // Enable gz compression by default
-            if (writeMode != WriteMode.TVF
-                    && streamLoadProp != null
-                    && !streamLoadProp.containsKey(COMPRESS_TYPE)) {
+            if (streamLoadProp != null && 
!streamLoadProp.containsKey(COMPRESS_TYPE)) {
                 streamLoadProp.put(COMPRESS_TYPE, COMPRESS_TYPE_GZ);
             }
 
+            if (writeMode == WriteMode.TVF && streamLoadProp != null) {
+                String compressType = 
streamLoadProp.getProperty(COMPRESS_TYPE).trim();
+                Preconditions.checkArgument(
+                        compressType.isEmpty() || 
COMPRESS_TYPE_GZ.equalsIgnoreCase(compressType),
+                        "TVF write mode only supports 'gz' or an empty 
compress_type.");
+            }
+
             Preconditions.checkArgument(
                     bufferFlushIntervalMs >= 1000,
                     "bufferFlushIntervalMs must be greater than or equal to 1 
second");
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/S3TvfOptions.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/S3TvfOptions.java
index b32a52e8..cbd63ae8 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/S3TvfOptions.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/cfg/S3TvfOptions.java
@@ -31,6 +31,8 @@ public class S3TvfOptions implements Serializable {
     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) {
@@ -40,6 +42,8 @@ public class S3TvfOptions implements Serializable {
         this.prefix = builder.prefix;
         this.accessKey = builder.accessKey;
         this.secretKey = builder.secretKey;
+        this.roleArn = builder.roleArn;
+        this.externalId = builder.externalId;
         this.pathStyleAccess = builder.pathStyleAccess;
     }
 
@@ -71,6 +75,22 @@ public class S3TvfOptions implements Serializable {
         return secretKey;
     }
 
+    public String getRoleArn() {
+        return roleArn;
+    }
+
+    public String getExternalId() {
+        return externalId;
+    }
+
+    public boolean hasRoleArn() {
+        return hasText(roleArn);
+    }
+
+    public boolean hasStaticCredentials() {
+        return hasText(accessKey);
+    }
+
     public boolean isPathStyleAccess() {
         return pathStyleAccess;
     }
@@ -90,13 +110,23 @@ public class S3TvfOptions implements Serializable {
                 && Objects.equals(bucket, that.bucket)
                 && Objects.equals(prefix, that.prefix)
                 && Objects.equals(accessKey, that.accessKey)
-                && Objects.equals(secretKey, that.secretKey);
+                && Objects.equals(secretKey, that.secretKey)
+                && Objects.equals(roleArn, that.roleArn)
+                && Objects.equals(externalId, that.externalId);
     }
 
     @Override
     public int hashCode() {
         return Objects.hash(
-                endpoint, region, bucket, prefix, accessKey, secretKey, 
pathStyleAccess);
+                endpoint,
+                region,
+                bucket,
+                prefix,
+                accessKey,
+                secretKey,
+                roleArn,
+                externalId,
+                pathStyleAccess);
     }
 
     @Override
@@ -127,6 +157,8 @@ public class S3TvfOptions implements Serializable {
         private String prefix;
         private String accessKey;
         private String secretKey;
+        private String roleArn;
+        private String externalId;
         private boolean pathStyleAccess;
 
         public Builder setEndpoint(String endpoint) {
@@ -159,6 +191,16 @@ public class S3TvfOptions implements Serializable {
             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;
@@ -173,7 +215,26 @@ public class S3TvfOptions implements Serializable {
                     }
                 }
             }
+            boolean hasAccessKey = hasText(accessKey);
+            boolean hasSecretKey = hasText(secretKey);
+            boolean hasRoleArn = hasText(roleArn);
+            if (hasAccessKey != hasSecretKey) {
+                throw new IllegalArgumentException(
+                        "sink.s3.access-key and sink.s3.secret-key must be 
configured together.");
+            }
+            if (!hasAccessKey && !hasRoleArn) {
+                throw new IllegalArgumentException(
+                        "S3 TVF options require either access/secret keys or 
sink.s3.role-arn.");
+            }
+            if (hasText(externalId) && !hasRoleArn) {
+                throw new IllegalArgumentException(
+                        "sink.s3.external-id requires sink.s3.role-arn.");
+            }
             return new S3TvfOptions(this);
         }
     }
+
+    private static boolean hasText(String value) {
+        return value != null && !value.trim().isEmpty();
+    }
 }
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 f8e25ff7..19ece5dc 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
@@ -19,6 +19,8 @@ package org.apache.doris.flink.sink.writer.tvf;
 
 import org.apache.doris.flink.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;
@@ -26,6 +28,9 @@ 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;
 
 import java.io.ByteArrayInputStream;
 import java.io.IOException;
@@ -35,12 +40,17 @@ import java.net.URI;
 public class S3ClientObjectStore implements S3ObjectStore {
 
     private static final String JSON_LINES_CONTENT_TYPE = 
"application/x-ndjson";
+    private static final String ROLE_SESSION_NAME = "doris-flink-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());
+        bucket = options.getBucket();
+        s3Client = createClient(options);
     }
 
     S3ClientObjectStore(S3Client s3Client, String bucket) {
@@ -48,14 +58,60 @@ public class S3ClientObjectStore implements S3ObjectStore {
         this.bucket = bucket;
     }
 
-    private static S3Client createClient(S3TvfOptions options) {
+    private AwsCredentialsProvider createCredentialsProvider(S3TvfOptions 
options) {
+        if (options.hasRoleArn()) {
+            return createAssumeRoleCredentialsProvider(options);
+        }
+        return staticCredentialsProvider(options);
+    }
+
+    private StsAssumeRoleCredentialsProvider 
createAssumeRoleCredentialsProvider(
+            S3TvfOptions options) {
+        AwsCredentialsProvider sourceCredentialsProvider =
+                createStsSourceCredentialsProvider(options);
+        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;
+    }
+
+    private AwsCredentialsProvider 
createStsSourceCredentialsProvider(S3TvfOptions options) {
+        if (options.hasStaticCredentials()) {
+            return staticCredentialsProvider(options);
+        }
+        defaultCredentialsProvider = 
DefaultCredentialsProvider.builder().build();
+        return defaultCredentialsProvider;
+    }
+
+    private static StaticCredentialsProvider 
staticCredentialsProvider(S3TvfOptions options) {
+        return StaticCredentialsProvider.create(
+                AwsBasicCredentials.create(options.getAccessKey(), 
options.getSecretKey()));
+    }
+
+    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 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()
@@ -89,6 +145,22 @@ public class S3ClientObjectStore implements S3ObjectStore {
 
     @Override
     public void close() {
-        s3Client.close();
+        try {
+            s3Client.close();
+        } finally {
+            closeCredentialsProviders();
+        }
+    }
+
+    private void closeCredentialsProviders() {
+        if (assumeRoleCredentialsProvider != null) {
+            assumeRoleCredentialsProvider.close();
+        }
+        if (stsClient != null) {
+            stsClient.close();
+        }
+        if (defaultCredentialsProvider != null) {
+            defaultCredentialsProvider.close();
+        }
     }
 }
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 4fcaf0ec..2089372e 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
@@ -34,6 +34,8 @@ import java.util.Map;
 import java.util.Properties;
 import java.util.concurrent.TimeUnit;
 
+import static org.apache.doris.flink.sink.writer.LoadConstants.COMPRESS_TYPE;
+
 /** Commits staged objects with one INSERT statement per writer and 
checkpoint. */
 public class S3TvfCommitter implements Committer<S3TvfCommittable> {
 
@@ -55,16 +57,18 @@ public class S3TvfCommitter implements 
Committer<S3TvfCommittable> {
                 new JdbcS3TvfLoadClient(dorisOptions),
                 executionOptions.getS3TvfOptions(),
                 executionOptions.getStreamLoadProp(),
-                executionOptions.getMaxRetries());
+                executionOptions.getMaxRetries(),
+                executionOptions.isGzipCompressionEnabled());
     }
 
     S3TvfCommitter(
             S3TvfLoadClient loadClient,
             S3TvfOptions options,
             Properties sessionProperties,
-            int maxRetries) {
+            int maxRetries,
+            boolean gzipEnabled) {
         this.loadClient = loadClient;
-        this.sqlBuilder = new S3TvfSqlBuilder(options);
+        this.sqlBuilder = new S3TvfSqlBuilder(options, gzipEnabled);
         this.sessionVariables = toSessionVariables(sessionProperties);
         this.maxRetries = maxRetries;
     }
@@ -202,7 +206,8 @@ public class S3TvfCommitter implements 
Committer<S3TvfCommittable> {
             if (!COLUMNS.equals(name)
                     && !PARTIAL_COLUMNS.equals(name)
                     && !FORMAT.equals(name)
-                    && !READ_JSON_BY_LINE.equals(name)) {
+                    && !READ_JSON_BY_LINE.equals(name)
+                    && !COMPRESS_TYPE.equals(name)) {
                 values.put(name, properties.getProperty(name));
             }
         }
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilder.java
 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilder.java
index cbe71573..fce43c03 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilder.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/main/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilder.java
@@ -33,9 +33,11 @@ import static 
org.apache.doris.flink.sink.writer.tvf.TvfSqlUtils.quoteLiteral;
 class S3TvfSqlBuilder {
 
     private final S3TvfOptions options;
+    private final boolean gzipEnabled;
 
-    public S3TvfSqlBuilder(S3TvfOptions options) {
+    public S3TvfSqlBuilder(S3TvfOptions options, boolean gzipEnabled) {
         this.options = options;
+        this.gzipEnabled = gzipEnabled;
     }
 
     public String buildInsertSql(S3TvfCommittable committable) {
@@ -46,6 +48,7 @@ class S3TvfSqlBuilder {
         }
         String columnSql = joinIdentifiers(loadColumns);
         String uri = buildUri(committable.getObjectKeys());
+        String credentials = buildCredentials();
 
         return "INSERT INTO "
                 + quoteIdentifier(committable.getDatabase())
@@ -60,9 +63,7 @@ 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())
                 + ","
@@ -71,6 +72,7 @@ class S3TvfSqlBuilder {
                 + property("format", "json")
                 + ","
                 + property("read_json_by_line", "true")
+                + (gzipEnabled ? "," + property("compress_type", "gz") : "")
                 + ","
                 + property("use_path_style", 
Boolean.toString(options.isPathStyleAccess()))
                 + ")";
@@ -83,6 +85,21 @@ class S3TvfSqlBuilder {
         return "s3://" + options.getBucket() + "/{" + String.join(",", 
objectKeys) + "}";
     }
 
+    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 static String joinIdentifiers(List<String> identifiers) {
         StringJoiner joiner = new StringJoiner(",");
         for (String identifier : identifiers) {
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 a53608ba..7454379b 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
@@ -38,6 +38,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;
 
 /** Shared writer that stages JSON Lines files in S3-compatible object 
storage. */
 public class S3TvfWriter<IN> {
@@ -56,6 +57,7 @@ public class S3TvfWriter<IN> {
     private final boolean deleteSignEnabled;
     private final int maxBytes;
     private final int uploadQueueSize;
+    private final boolean gzipEnabled;
     private final ByteArrayOutputStream buffer = new ByteArrayOutputStream();
     private final List<String> currentObjectKeys = new ArrayList<>();
     private final BlockingQueue<Runnable> uploadQueue;
@@ -77,7 +79,8 @@ public class S3TvfWriter<IN> {
             List<String> columns,
             boolean deleteSignEnabled,
             int maxBytes,
-            int uploadQueueSize) {
+            int uploadQueueSize,
+            boolean gzipEnabled) {
         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;
@@ -92,6 +95,7 @@ public class S3TvfWriter<IN> {
         this.deleteSignEnabled = deleteSignEnabled;
         this.maxBytes = maxBytes;
         this.uploadQueueSize = uploadQueueSize;
+        this.gzipEnabled = gzipEnabled;
         this.uploadQueue = new LinkedBlockingQueue<>(uploadQueueSize);
         this.uploadExecutor =
                 Executors.newSingleThreadExecutor(
@@ -173,8 +177,13 @@ public class S3TvfWriter<IN> {
         }
         String fileName =
                 String.format(
-                        "%s_%s_%d_%d_%d.json",
-                        labelPrefix, table, subtaskId, currentCheckpointId, 
fileNumber++);
+                        "%s_%s_%d_%d_%d.json%s",
+                        labelPrefix,
+                        table,
+                        subtaskId,
+                        currentCheckpointId,
+                        fileNumber++,
+                        gzipEnabled ? ".gz" : "");
         String objectKey = objectPrefix + (objectPrefix.endsWith("/") ? "" : 
"/") + fileName;
         byte[] content = buffer.toByteArray();
         putUpload(
@@ -184,14 +193,15 @@ public class S3TvfWriter<IN> {
                     }
                     long uploadStartedAtNanos = System.nanoTime();
                     try {
-                        objectStore.put(objectKey, content);
+                        byte[] uploadContent = gzipEnabled ? gzip(content) : 
content;
+                        objectStore.put(objectKey, uploadContent);
                         currentObjectKeys.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) {
@@ -215,6 +225,14 @@ public class S3TvfWriter<IN> {
         buffer.reset();
     }
 
+    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 processUploads() {
         while (!Thread.currentThread().isInterrupted()) {
             try {
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 e83df655..3b159679 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
@@ -44,6 +44,8 @@ import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_REQUEST_RETR
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_TABLET_SIZE_DEFAULT;
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.DORIS_THRIFT_MAX_MESSAGE_SIZE_DEFAULT;
 import static 
org.apache.doris.flink.cfg.ConfigurationOptions.SOURCE_BINLOG_VISIBLE_WAIT_TIMEOUT_MS_DEFAULT;
+import static org.apache.doris.flink.sink.writer.LoadConstants.COMPRESS_TYPE;
+import static 
org.apache.doris.flink.sink.writer.LoadConstants.COMPRESS_TYPE_GZ;
 import static org.apache.doris.flink.sink.writer.LoadConstants.FORMAT_KEY;
 import static org.apache.doris.flink.sink.writer.LoadConstants.JSON;
 import static 
org.apache.doris.flink.sink.writer.LoadConstants.READ_JSON_BY_LINE;
@@ -328,6 +330,18 @@ public class DorisConfigOptions {
                     .noDefaultValue()
                     .withDescription("Secret key of the S3-compatible object 
storage.");
 
+    public static final ConfigOption<String> SINK_S3_ROLE_ARN =
+            ConfigOptions.key("sink.s3.role-arn")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("AWS IAM role ARN used to access the S3 
object storage.");
+
+    public static final ConfigOption<String> SINK_S3_EXTERNAL_ID =
+            ConfigOptions.key("sink.s3.external-id")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("External ID used when assuming the AWS 
IAM role.");
+
     public static final ConfigOption<Boolean> SINK_S3_PATH_STYLE_ACCESS =
             ConfigOptions.key("sink.s3.path-style-access")
                     .booleanType()
@@ -472,8 +486,22 @@ public class DorisConfigOptions {
         String region = requireNonBlank(readableConfig, SINK_S3_REGION);
         String bucket = requireNonBlank(readableConfig, SINK_S3_BUCKET);
         String prefix = requireNonBlank(readableConfig, SINK_S3_PREFIX);
-        String accessKey = requireNonBlank(readableConfig, SINK_S3_ACCESS_KEY);
-        String secretKey = requireNonBlank(readableConfig, SINK_S3_SECRET_KEY);
+        String accessKey = optionalNonBlank(readableConfig, 
SINK_S3_ACCESS_KEY);
+        String secretKey = optionalNonBlank(readableConfig, 
SINK_S3_SECRET_KEY);
+        String roleArn = optionalNonBlank(readableConfig, SINK_S3_ROLE_ARN);
+        String externalId = optionalNonBlank(readableConfig, 
SINK_S3_EXTERNAL_ID);
+        if ((accessKey == null) != (secretKey == null)) {
+            throw new ValidationException(
+                    "Options 'sink.s3.access-key' and 'sink.s3.secret-key' 
must be configured together.");
+        }
+        if (roleArn == null && accessKey == null) {
+            throw new ValidationException(
+                    "TVF write mode requires either S3 access/secret keys or 
'sink.s3.role-arn'.");
+        }
+        if (externalId != null && roleArn == null) {
+            throw new ValidationException(
+                    "Option 'sink.s3.external-id' requires 
'sink.s3.role-arn'.");
+        }
 
         return S3TvfOptions.builder()
                 .setEndpoint(endpoint)
@@ -482,6 +510,8 @@ public class DorisConfigOptions {
                 .setPrefix(prefix)
                 .setAccessKey(accessKey)
                 .setSecretKey(secretKey)
+                .setRoleArn(roleArn)
+                .setExternalId(externalId)
                 
.setPathStyleAccess(readableConfig.get(SINK_S3_PATH_STYLE_ACCESS))
                 .build();
     }
@@ -496,16 +526,31 @@ public class DorisConfigOptions {
             throw new ValidationException(
                     "TVF write mode requires 
'sink.properties.read_json_by_line' to be true.");
         }
+        validateTvfCompression(loadProperties);
+    }
+
+    private static void validateTvfCompression(Properties loadProperties) {
+        String compressType = loadProperties.getProperty(COMPRESS_TYPE, 
COMPRESS_TYPE_GZ).trim();
+        if (!compressType.isEmpty() && 
!COMPRESS_TYPE_GZ.equalsIgnoreCase(compressType)) {
+            throw new ValidationException(
+                    "TVF write mode only supports 'gz' or an empty 
compress_type.");
+        }
     }
 
     private static String requireNonBlank(
             ReadableConfig readableConfig, ConfigOption<String> option) {
-        String value = readableConfig.getOptional(option).orElse(null);
-        if (value == null || value.trim().isEmpty()) {
+        String value = optionalNonBlank(readableConfig, option);
+        if (value == null) {
             throw new ValidationException(
                     String.format("Option '%s' is required for TVF write 
mode.", option.key()));
         }
-        return value.trim();
+        return value;
+    }
+
+    private static String optionalNonBlank(
+            ReadableConfig readableConfig, ConfigOption<String> option) {
+        String value = readableConfig.getOptional(option).orElse(null);
+        return value == null || value.trim().isEmpty() ? null : value.trim();
     }
 
     public static final ConfigOption<Boolean> SINK_HTTP_UTF8_CHARSET =
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/DorisExecutionOptionsTest.java
 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/DorisExecutionOptionsTest.java
index a5fdd05c..69aa6360 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/DorisExecutionOptionsTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/DorisExecutionOptionsTest.java
@@ -192,7 +192,7 @@ public class DorisExecutionOptionsTest {
     }
 
     @Test
-    public void 
testTvfPropertiesAreSessionVariablesWithoutStreamLoadDefaults() {
+    public void testTvfPropertiesIncludeDefaultCompression() {
         Properties sessionVariables = new Properties();
         sessionVariables.setProperty("enable_unique_key_partial_update", 
"true");
 
@@ -205,7 +205,39 @@ public class DorisExecutionOptionsTest {
                         .build();
 
         Assert.assertEquals(sessionVariables, 
executionOptions.getStreamLoadProp());
-        Assert.assertEquals(1, executionOptions.getStreamLoadProp().size());
+        Assert.assertEquals(
+                "gz", 
executionOptions.getStreamLoadProp().getProperty("compress_type"));
+        Assert.assertTrue(executionOptions.isGzipCompressionEnabled());
+        Assert.assertEquals(2, executionOptions.getStreamLoadProp().size());
+    }
+
+    @Test
+    public void testTvfAllowsDisablingCompression() {
+        Properties properties = new Properties();
+        properties.setProperty("compress_type", "");
+
+        DorisExecutionOptions executionOptions =
+                DorisExecutionOptions.builder()
+                        .setWriteMode(WriteMode.TVF)
+                        .setLabelPrefix("label")
+                        .setStreamLoadProp(properties)
+                        .setS3TvfOptions(s3TvfOptions())
+                        .build();
+
+        Assert.assertFalse(executionOptions.isGzipCompressionEnabled());
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testTvfRejectsUnsupportedCompression() {
+        Properties properties = new Properties();
+        properties.setProperty("compress_type", "zstd");
+
+        DorisExecutionOptions.builder()
+                .setWriteMode(WriteMode.TVF)
+                .setLabelPrefix("label")
+                .setStreamLoadProp(properties)
+                .setS3TvfOptions(s3TvfOptions())
+                .build();
     }
 
     @Test(expected = IllegalArgumentException.class)
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/S3TvfOptionsTest.java
 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/S3TvfOptionsTest.java
index 0b5b9747..bf1d6257 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/S3TvfOptionsTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/cfg/S3TvfOptionsTest.java
@@ -32,6 +32,8 @@ public class S3TvfOptionsTest {
                         .setPrefix("doris")
                         .setAccessKey("access-key")
                         .setSecretKey("secret-key")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
                         .setPathStyleAccess(true)
                         .build();
 
@@ -41,20 +43,48 @@ public class S3TvfOptionsTest {
         Assert.assertEquals("doris", 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
     public void testRejectsGlobCharactersInPrefix() {
         for (String character : new String[] {"*", "?", "[", "]", "{", "}", 
",", "\\"}) {
             try {
-                S3TvfOptions.builder().setPrefix("path/" + character + 
"/prefix").build();
+                S3TvfOptions.builder()
+                        .setPrefix("path/" + character + "/prefix")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .build();
                 Assert.fail("Expected prefix containing '" + character + "' to 
be rejected.");
             } catch (IllegalArgumentException expected) {
                 // Expected.
             }
         }
     }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testRequiresCredentialsOrRole() {
+        S3TvfOptions.builder().build();
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testRequiresCompleteStaticCredentials() {
+        S3TvfOptions.builder()
+                .setAccessKey("access-key")
+                .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                .build();
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testExternalIdRequiresRole() {
+        S3TvfOptions.builder()
+                .setAccessKey("access-key")
+                .setSecretKey("secret-key")
+                .setExternalId("external-id")
+                .build();
+    }
 }
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 a3c83f14..4298fcaa 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
@@ -17,12 +17,14 @@
 
 package org.apache.doris.flink.sink.writer.tvf;
 
+import org.apache.doris.flink.cfg.S3TvfOptions;
 import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.ArgumentCaptor;
 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.sts.model.AssumeRoleRequest;
 
 import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
@@ -32,6 +34,21 @@ import static org.mockito.Mockito.verify;
 
 public class S3ClientObjectStoreTest {
 
+    @Test
+    public void testBuildAssumeRoleRequest() {
+        S3TvfOptions options =
+                S3TvfOptions.builder()
+                        .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-flink-connector", 
request.roleSessionName());
+    }
+
     @Test
     public void testPutObjectUsesRepeatableContentProviderWithoutCopying() 
throws Exception {
         S3Client s3Client = mock(S3Client.class);
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 8168d8cf..70b2f322 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
@@ -81,17 +81,31 @@ public class S3TvfCommitterTest {
         properties.setProperty("columns", "id,name");
         properties.setProperty("partial_columns", "true");
         properties.setProperty("query_timeout", "60");
-        S3TvfCommitter committer = new S3TvfCommitter(loadClient, options(), 
properties, 0);
+        properties.setProperty("compress_type", "gz");
+        S3TvfCommitter committer = new S3TvfCommitter(loadClient, options(), 
properties, 0, true);
 
         
committer.commit(Collections.singletonList(request(committable("file.json"))));
 
         Assert.assertFalse(loadClient.sessionVariables.containsKey("columns"));
         
Assert.assertFalse(loadClient.sessionVariables.containsKey("partial_columns"));
+        
Assert.assertFalse(loadClient.sessionVariables.containsKey("compress_type"));
         Assert.assertEquals(
                 "true", 
loadClient.sessionVariables.get("enable_unique_key_partial_update"));
         Assert.assertEquals("60", 
loadClient.sessionVariables.get("query_timeout"));
     }
 
+    @Test
+    public void testEmptyCompressTypeDisablesGzipInInsertSql() throws 
Exception {
+        RecordingLoadClient loadClient = new RecordingLoadClient();
+        Properties properties = new Properties();
+        properties.setProperty("compress_type", "");
+        S3TvfCommitter committer = new S3TvfCommitter(loadClient, options(), 
properties, 0, false);
+
+        
committer.commit(Collections.singletonList(request(committable("file.json"))));
+
+        Assert.assertFalse(loadClient.lastInsertSql.contains("'compress_type' 
= 'gz'"));
+    }
+
     @Test
     public void testLabelAlreadyUsedAndFinishedIsSuccessful() throws Exception 
{
         RecordingLoadClient loadClient = new RecordingLoadClient();
@@ -202,7 +216,7 @@ public class S3TvfCommitterTest {
     private static S3TvfCommitter createCommitter(RecordingLoadClient 
loadClient, int maxRetries) {
         Properties sessionVariables = new Properties();
         sessionVariables.setProperty("enable_partial_update", "true");
-        return new S3TvfCommitter(loadClient, options(), sessionVariables, 
maxRetries);
+        return new S3TvfCommitter(loadClient, options(), sessionVariables, 
maxRetries, true);
     }
 
     private static S3TvfOptions options() {
diff --git 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilderTest.java
 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilderTest.java
index 2c0b83dc..1b4d849e 100644
--- 
a/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilderTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-base/src/test/java/org/apache/doris/flink/sink/writer/tvf/S3TvfSqlBuilderTest.java
@@ -47,7 +47,7 @@ public class S3TvfSqlBuilderTest {
                         Arrays.asList("id", "name"),
                         true);
 
-        String sql = new S3TvfSqlBuilder(options).buildInsertSql(committable);
+        String sql = new S3TvfSqlBuilder(options, 
true).buildInsertSql(committable);
 
         Assert.assertEquals(
                 "INSERT INTO `db`.`tbl` WITH LABEL `label_tbl_7` "
@@ -57,7 +57,58 @@ public class S3TvfSqlBuilderTest {
                         + "'s3.access_key' = 'ak','s3.secret_key' = 'sk',"
                         + "'s3.region' = 'us-east-1','s3.endpoint' = 
'https://s3.example.com',"
                         + "'format' = 'json','read_json_by_line' = 'true',"
+                        + "'compress_type' = 'gz',"
                         + "'use_path_style' = 'true')",
                 sql);
     }
+
+    @Test
+    public void testBuildInsertSqlWithIamRole() {
+        S3TvfOptions options =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.us-east-1.amazonaws.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("bucket")
+                        .setPrefix("prefix")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
+                        .build();
+        S3TvfCommittable committable =
+                new S3TvfCommittable(
+                        7L,
+                        "db",
+                        "tbl",
+                        "label_tbl_7",
+                        Arrays.asList("prefix_tbl_0_7_0.json"),
+                        Arrays.asList("id"),
+                        false);
+
+        String sql = new S3TvfSqlBuilder(options, 
true).buildInsertSql(committable);
+
+        Assert.assertTrue(
+                sql.contains(
+                        "'s3.role_arn' = 
'arn:aws:iam::123456789012:role/doris',"
+                                + "'s3.external_id' = 'external-id'"));
+        Assert.assertFalse(sql.contains("s3.access_key"));
+        Assert.assertFalse(sql.contains("s3.secret_key"));
+        Assert.assertTrue(sql.contains("'compress_type' = 'gz'"));
+
+        S3TvfOptions optionsWithSourceCredentials =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.us-east-1.amazonaws.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("bucket")
+                        .setPrefix("prefix")
+                        .setAccessKey("ak")
+                        .setSecretKey("sk")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .build();
+        String sqlWithSourceCredentials =
+                new S3TvfSqlBuilder(optionsWithSourceCredentials, 
true).buildInsertSql(committable);
+        Assert.assertTrue(sqlWithSourceCredentials.contains("'s3.access_key' = 
'ak'"));
+        Assert.assertTrue(sqlWithSourceCredentials.contains("'s3.secret_key' = 
'sk'"));
+        Assert.assertTrue(
+                sqlWithSourceCredentials.contains(
+                        "'s3.role_arn' = 
'arn:aws:iam::123456789012:role/doris'"));
+    }
 }
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 565e123b..c9d122e5 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
@@ -26,6 +26,8 @@ import org.junit.Rule;
 import org.junit.Test;
 import org.slf4j.event.Level;
 
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
@@ -38,6 +40,7 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
+import java.util.zip.GZIPInputStream;
 
 public class S3TvfWriterTest {
 
@@ -107,6 +110,19 @@ public class S3TvfWriterTest {
         Assert.assertEquals("label_tbl_2_8", 
writer.prepareCommit().iterator().next().getLabel());
     }
 
+    @Test
+    public void testGzipUpload() throws Exception {
+        RecordingObjectStore objectStore = new RecordingObjectStore();
+        S3TvfWriter<String> writer = createWriter(6L, objectStore, 2, 10, 
true);
+
+        writer.write("12345");
+        writer.flush();
+
+        Assert.assertEquals("prefix/label_tbl_2_7_0.json.gz", 
objectStore.objectKeys.get(0));
+        Assert.assertArrayEquals(
+                "12345\n".getBytes(StandardCharsets.UTF_8), 
gunzip(objectStore.contents.get(0)));
+    }
+
     @Test
     public void testAdvanceAcrossSkippedCheckpointId() throws Exception {
         RecordingObjectStore objectStore = new RecordingObjectStore();
@@ -192,7 +208,7 @@ public class S3TvfWriterTest {
 
     private static S3TvfWriter<String> createWriter(
             long restoredCheckpointId, RecordingObjectStore objectStore) {
-        return createWriter(restoredCheckpointId, objectStore, 2, 10);
+        return createWriter(restoredCheckpointId, objectStore, 2, 10, false);
     }
 
     private static S3TvfWriter<String> createWriter(
@@ -200,6 +216,15 @@ public class S3TvfWriterTest {
             RecordingObjectStore objectStore,
             int uploadQueueSize,
             int maxBytes) {
+        return createWriter(restoredCheckpointId, objectStore, 
uploadQueueSize, maxBytes, false);
+    }
+
+    private static S3TvfWriter<String> createWriter(
+            long restoredCheckpointId,
+            RecordingObjectStore objectStore,
+            int uploadQueueSize,
+            int maxBytes,
+            boolean gzipEnabled) {
         DorisRecordSerializer<String> serializer =
                 value -> 
DorisRecord.of(value.getBytes(StandardCharsets.UTF_8));
         return new S3TvfWriter<>(
@@ -214,7 +239,20 @@ public class S3TvfWriterTest {
                 Arrays.asList("id", "name"),
                 true,
                 maxBytes,
-                uploadQueueSize);
+                uploadQueueSize,
+                gzipEnabled);
+    }
+
+    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 {
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 7cb2d8b1..8bbf4126 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
@@ -75,7 +75,8 @@ public class S3TvfWriterAdapter<IN>
                         rowDataSerializer.getSelectedColumns(),
                         rowDataSerializer.isDeleteSignEnabled(),
                         executionOptions.getBufferFlushMaxBytes(),
-                        executionOptions.getFlushQueueSize());
+                        executionOptions.getFlushQueueSize(),
+                        executionOptions.isGzipCompressionEnabled());
     }
 
     @Override
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
index 34723646..ac053468 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
@@ -89,9 +89,11 @@ import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_PARALLELISM;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ACCESS_KEY;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_BUCKET;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ENDPOINT;
+import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_EXTERNAL_ID;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_PATH_STYLE_ACCESS;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_PREFIX;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_REGION;
+import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ROLE_ARN;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_SECRET_KEY;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_USE_CACHE;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_WRITE_MODE;
@@ -204,6 +206,8 @@ public final class DorisDynamicTableFactory
         options.add(SINK_S3_PREFIX);
         options.add(SINK_S3_ACCESS_KEY);
         options.add(SINK_S3_SECRET_KEY);
+        options.add(SINK_S3_ROLE_ARN);
+        options.add(SINK_S3_EXTERNAL_ID);
         options.add(SINK_S3_PATH_STYLE_ACCESS);
         return options;
     }
diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
index 95b0e54e..298c6805 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
@@ -289,8 +289,8 @@ public class DorisDynamicTableFactoryTest {
         properties.put("sink.s3.region", "us-east-1");
         properties.put("sink.s3.bucket", "bucket");
         properties.put("sink.s3.prefix", "prefix");
-        properties.put("sink.s3.access-key", "ak");
-        properties.put("sink.s3.secret-key", "sk");
+        properties.put("sink.s3.role-arn", 
"arn:aws:iam::123456789012:role/doris");
+        properties.put("sink.s3.external-id", "external-id");
         properties.put("sink.s3.path-style-access", "true");
         properties.put("sink.properties.columns", "a,c");
 
@@ -314,8 +314,8 @@ public class DorisDynamicTableFactoryTest {
                         .setRegion("us-east-1")
                         .setBucket("bucket")
                         .setPrefix("prefix")
-                        .setAccessKey("ak")
-                        .setSecretKey("sk")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
                         .setPathStyleAccess(true)
                         .build();
         DorisExecutionOptions executionOptions =
@@ -340,6 +340,44 @@ public class DorisDynamicTableFactoryTest {
                         TableSchema.fromResolvedSchema(SCHEMA),
                         null);
         assertEquals(expected, actual);
+
+        properties.remove("sink.s3.role-arn");
+        properties.remove("sink.s3.external-id");
+        properties.put("sink.s3.access-key", "ak");
+        properties.put("sink.s3.secret-key", "sk");
+        actual = (DorisDynamicTableSink) FactoryMocks.createTableSink(SCHEMA, 
properties);
+
+        s3TvfOptions =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.example.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("bucket")
+                        .setPrefix("prefix")
+                        .setAccessKey("ak")
+                        .setSecretKey("sk")
+                        .setPathStyleAccess(true)
+                        .build();
+        executionOptions =
+                DorisExecutionOptions.builder()
+                        .setWriteMode(WriteMode.TVF)
+                        .setLabelPrefix("flink")
+                        .setBufferFlushMaxBytes(10 * 1024 * 1024)
+                        .setStreamLoadProp(
+                                new Properties() {
+                                    {
+                                        setProperty("columns", "a,c");
+                                    }
+                                })
+                        .setS3TvfOptions(s3TvfOptions)
+                        .build();
+        expected =
+                new DorisDynamicTableSink(
+                        options,
+                        DorisReadOptions.builder().build(),
+                        executionOptions,
+                        TableSchema.fromResolvedSchema(SCHEMA),
+                        null);
+        assertEquals(expected, actual);
     }
 
     private Map<String, String> getAllOptions() {
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 718eb147..ac53713d 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
@@ -75,7 +75,8 @@ public class S3TvfWriterAdapter<IN>
                         rowDataSerializer.getSelectedColumns(),
                         rowDataSerializer.isDeleteSignEnabled(),
                         executionOptions.getBufferFlushMaxBytes(),
-                        executionOptions.getFlushQueueSize());
+                        executionOptions.getFlushQueueSize(),
+                        executionOptions.isGzipCompressionEnabled());
     }
 
     @Override
diff --git 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
index 72461062..e730923a 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink2/src/main/java/org/apache/doris/flink/table/DorisDynamicTableFactory.java
@@ -89,9 +89,11 @@ import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_PARALLELISM;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ACCESS_KEY;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_BUCKET;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ENDPOINT;
+import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_EXTERNAL_ID;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_PATH_STYLE_ACCESS;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_PREFIX;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_REGION;
+import static org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_ROLE_ARN;
 import static 
org.apache.doris.flink.table.DorisConfigOptions.SINK_S3_SECRET_KEY;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_USE_CACHE;
 import static org.apache.doris.flink.table.DorisConfigOptions.SINK_WRITE_MODE;
@@ -204,6 +206,8 @@ public final class DorisDynamicTableFactory
         options.add(SINK_S3_PREFIX);
         options.add(SINK_S3_ACCESS_KEY);
         options.add(SINK_S3_SECRET_KEY);
+        options.add(SINK_S3_ROLE_ARN);
+        options.add(SINK_S3_EXTERNAL_ID);
         options.add(SINK_S3_PATH_STYLE_ACCESS);
         return options;
     }
diff --git 
a/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
 
b/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
index 754c33fe..38a2b8fc 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink2/src/test/java/org/apache/doris/flink/table/DorisDynamicTableFactoryTest.java
@@ -293,8 +293,8 @@ public class DorisDynamicTableFactoryTest {
         properties.put("sink.s3.region", "us-east-1");
         properties.put("sink.s3.bucket", "bucket");
         properties.put("sink.s3.prefix", "prefix");
-        properties.put("sink.s3.access-key", "ak");
-        properties.put("sink.s3.secret-key", "sk");
+        properties.put("sink.s3.role-arn", 
"arn:aws:iam::123456789012:role/doris");
+        properties.put("sink.s3.external-id", "external-id");
         properties.put("sink.s3.path-style-access", "true");
         properties.put("sink.properties.columns", "a,c");
 
@@ -318,8 +318,8 @@ public class DorisDynamicTableFactoryTest {
                         .setRegion("us-east-1")
                         .setBucket("bucket")
                         .setPrefix("prefix")
-                        .setAccessKey("ak")
-                        .setSecretKey("sk")
+                        .setRoleArn("arn:aws:iam::123456789012:role/doris")
+                        .setExternalId("external-id")
                         .setPathStyleAccess(true)
                         .build();
         DorisExecutionOptions executionOptions =
@@ -344,6 +344,44 @@ public class DorisDynamicTableFactoryTest {
                         TableSchema.fromResolvedSchema(SCHEMA),
                         null);
         assertEquals(expected, actual);
+
+        properties.remove("sink.s3.role-arn");
+        properties.remove("sink.s3.external-id");
+        properties.put("sink.s3.access-key", "ak");
+        properties.put("sink.s3.secret-key", "sk");
+        actual = (DorisDynamicTableSink) FactoryMocks.createTableSink(SCHEMA, 
properties);
+
+        s3TvfOptions =
+                S3TvfOptions.builder()
+                        .setEndpoint("https://s3.example.com";)
+                        .setRegion("us-east-1")
+                        .setBucket("bucket")
+                        .setPrefix("prefix")
+                        .setAccessKey("ak")
+                        .setSecretKey("sk")
+                        .setPathStyleAccess(true)
+                        .build();
+        executionOptions =
+                DorisExecutionOptions.builder()
+                        .setWriteMode(WriteMode.TVF)
+                        .setLabelPrefix("flink")
+                        .setBufferFlushMaxBytes(10 * 1024 * 1024)
+                        .setStreamLoadProp(
+                                new Properties() {
+                                    {
+                                        setProperty("columns", "a,c");
+                                    }
+                                })
+                        .setS3TvfOptions(s3TvfOptions)
+                        .build();
+        expected =
+                new DorisDynamicTableSink(
+                        options,
+                        DorisReadOptions.builder().build(),
+                        executionOptions,
+                        TableSchema.fromResolvedSchema(SCHEMA),
+                        null);
+        assertEquals(expected, actual);
     }
 
     private Map<String, String> getAllOptions() {
diff --git 
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfIamRoleITCase.java
 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfIamRoleITCase.java
new file mode 100644
index 00000000..369378df
--- /dev/null
+++ 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfIamRoleITCase.java
@@ -0,0 +1,203 @@
+// 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.flink.sink;
+
+import org.apache.flink.api.common.RuntimeExecutionMode;
+import org.apache.flink.runtime.minicluster.RpcServiceSharing;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.StringData;
+import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.test.util.MiniClusterWithClientResource;
+
+import org.apache.doris.flink.cfg.DorisExecutionOptions;
+import org.apache.doris.flink.cfg.DorisOptions;
+import org.apache.doris.flink.cfg.DorisReadOptions;
+import org.apache.doris.flink.cfg.S3TvfOptions;
+import org.apache.doris.flink.container.ContainerUtils;
+import org.apache.doris.flink.container.instance.DorisCustomerContainer;
+import org.apache.doris.flink.sink.writer.WriteMode;
+import org.apache.doris.flink.sink.writer.tvf.S3TvfRowDataSerializer;
+import org.junit.AfterClass;
+import org.junit.Assume;
+import org.junit.BeforeClass;
+import org.junit.Rule;
+import org.junit.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.sql.Connection;
+import java.util.Arrays;
+import java.util.UUID;
+
+/** Opt-in integration test for S3 TVF writes with an AWS IAM role. */
+public class S3TvfIamRoleITCase {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(S3TvfIamRoleITCase.class);
+    private static final String DATABASE = "test_s3_tvf_iam_role";
+    private static DorisCustomerContainer doris;
+
+    @Rule
+    public final MiniClusterWithClientResource miniClusterResource =
+            new MiniClusterWithClientResource(
+                    new MiniClusterResourceConfiguration.Builder()
+                            .setNumberTaskManagers(1)
+                            .setNumberSlotsPerTaskManager(1)
+                            .setRpcServiceSharing(RpcServiceSharing.DEDICATED)
+                            .build());
+
+    @BeforeClass
+    public static void useExternalEnvironment() {
+        Assume.assumeTrue(
+                "IAM role ITCase requires -Ds3_tvf_iam_role_it=true",
+                Boolean.getBoolean("s3_tvf_iam_role_it"));
+        Assume.assumeTrue(
+                "IAM role ITCase requires -Dcustomer_env=true", 
Boolean.getBoolean("customer_env"));
+        requiredProperty("s3_endpoint");
+        requiredProperty("s3_region");
+        requiredProperty("s3_bucket");
+        requiredProperty("s3_role_arn");
+
+        doris = new DorisCustomerContainer();
+        doris.startContainer();
+    }
+
+    @AfterClass
+    public static void closeExternalEnvironment() {
+        if (doris != null) {
+            doris.close();
+        }
+    }
+
+    @Test
+    public void testWritesThroughIamRole() throws Exception {
+        String table = "iam_role_" + UUID.randomUUID().toString().replace("-", 
"");
+        try {
+            createTable(table);
+            runTvfSink(table);
+            assertRows(table);
+        } finally {
+            dropTable(table);
+        }
+    }
+
+    private void runTvfSink(String table) throws Exception {
+        String[] fieldNames = {"id", "name"};
+        DataType[] dataTypes = {DataTypes.INT(), DataTypes.STRING()};
+        LogicalType[] logicalTypes =
+                
Arrays.stream(dataTypes).map(DataType::getLogicalType).toArray(LogicalType[]::new);
+        InternalTypeInfo<RowData> typeInfo = 
InternalTypeInfo.ofFields(logicalTypes, fieldNames);
+
+        S3TvfOptions s3Options =
+                S3TvfOptions.builder()
+                        .setEndpoint(requiredProperty("s3_endpoint"))
+                        .setRegion(requiredProperty("s3_region"))
+                        .setBucket(requiredProperty("s3_bucket"))
+                        .setPrefix(System.getProperty("s3_prefix", 
"doris-flink-connector-it"))
+                        .setRoleArn(requiredProperty("s3_role_arn"))
+                        .setExternalId(optionalProperty("s3_external_id"))
+                        
.setPathStyleAccess(Boolean.getBoolean("s3_path_style_access"))
+                        .build();
+        DorisSink<RowData> sink =
+                DorisSink.<RowData>builder()
+                        .setDorisOptions(
+                                DorisOptions.builder()
+                                        .setFenodes(doris.getFenodes())
+                                        .setJdbcUrl(doris.getJdbcUrl())
+                                        .setTableIdentifier(DATABASE + "." + 
table)
+                                        .setUsername(doris.getUsername())
+                                        .setPassword(doris.getPassword())
+                                        .build())
+                        
.setDorisReadOptions(DorisReadOptions.builder().build())
+                        .setDorisExecutionOptions(
+                                DorisExecutionOptions.builder()
+                                        .setWriteMode(WriteMode.TVF)
+                                        .setLabelPrefix("iam_role_" + 
UUID.randomUUID())
+                                        .setBufferFlushMaxBytes(1024)
+                                        .setS3TvfOptions(s3Options)
+                                        .build())
+                        .setSerializer(
+                                new S3TvfRowDataSerializer(
+                                        fieldNames, dataTypes, 
Arrays.asList(fieldNames), false))
+                        .build();
+
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+        env.setRuntimeMode(RuntimeExecutionMode.BATCH);
+        env.setParallelism(1);
+        env.fromCollection(Arrays.asList(row(1, "doris"), row(2, "flink")), 
typeInfo)
+                .sinkTo(sink)
+                .setParallelism(1);
+        env.execute("S3 TVF IAM role ITCase");
+    }
+
+    private static void createTable(String table) {
+        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\")");
+    }
+
+    private static void assertRows(String table) {
+        ContainerUtils.checkResult(
+                doris.getQueryConnection(),
+                LOG,
+                Arrays.asList("1,doris", "2,flink"),
+                "SELECT id,name FROM `" + DATABASE + "`.`" + table + "` ORDER 
BY id",
+                2,
+                true);
+    }
+
+    private static void dropTable(String table) {
+        executeSql("DROP TABLE IF EXISTS `" + DATABASE + "`.`" + table + "`");
+    }
+
+    private static void executeSql(String... statements) {
+        try (Connection connection = doris.getQueryConnection()) {
+            ContainerUtils.executeSQLStatement(connection, LOG, statements);
+        } catch (Exception e) {
+            throw new RuntimeException(e);
+        }
+    }
+
+    private static GenericRowData row(int id, String name) {
+        return GenericRowData.of(id, StringData.fromString(name));
+    }
+
+    private static String requiredProperty(String name) {
+        String value = optionalProperty(name);
+        if (value == null) {
+            throw new IllegalArgumentException("Missing required system 
property: " + name);
+        }
+        return value;
+    }
+
+    private static String optionalProperty(String name) {
+        String value = System.getProperty(name);
+        return value == null || value.trim().isEmpty() ? null : value.trim();
+    }
+}
diff --git 
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
index 97de848b..58858404 100644
--- 
a/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
+++ 
b/flink-doris-connector/flink-doris-connector-it/src/test/java/org/apache/doris/flink/sink/S3TvfSinkITCase.java
@@ -68,6 +68,7 @@ import software.amazon.awssdk.services.s3.S3Client;
 import software.amazon.awssdk.services.s3.S3Configuration;
 import software.amazon.awssdk.services.s3.model.CreateBucketRequest;
 import software.amazon.awssdk.services.s3.model.ListObjectsV2Request;
+import software.amazon.awssdk.services.s3.model.ListObjectsV2Response;
 
 import java.net.Inet4Address;
 import java.net.InetAddress;
@@ -198,14 +199,80 @@ public class S3TvfSinkITCase extends 
AbstractITCaseService {
                 "id,name,note",
                 Arrays.asList("1,doris,中文", "2,flink,quote-'and-\"", 
"3,null-value,null"),
                 3);
-        int objectCount =
+        ListObjectsV2Response objects =
                 s3Client.listObjectsV2(
-                                ListObjectsV2Request.builder()
-                                        .bucket(BUCKET)
-                                        .prefix(objectPrefix + "/" + 
labelPrefix + "_" + table)
-                                        .build())
-                        .keyCount();
-        Assert.assertTrue("Expected the buffer limit to create multiple 
objects", objectCount > 1);
+                        ListObjectsV2Request.builder()
+                                .bucket(BUCKET)
+                                .prefix(objectPrefix + "/" + labelPrefix + "_" 
+ table)
+                                .build());
+        Assert.assertTrue(
+                "Expected the buffer limit to create multiple objects", 
objects.keyCount() > 1);
+        Assert.assertTrue(
+                "Expected TVF objects to use gzip compression by default",
+                objects.contents().stream().allMatch(object -> 
object.key().endsWith(".json.gz")));
+    }
+
+    @Test
+    public void testDisablesCompressionWithEmptyCompressType() throws 
Exception {
+        String table = uniqueName("uncompressed");
+        String objectPrefix = uniqueName("uncompressed_objects");
+        String labelPrefix = uniqueName("uncompressed_label");
+        createDuplicateTable(table, "`id` INT, `name` VARCHAR(128)");
+
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+        env.setRuntimeMode(RuntimeExecutionMode.BATCH);
+        env.setParallelism(1);
+        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
+        tableEnv.executeSql(
+                String.format(
+                        "CREATE TABLE tvf_sink (id INT, name STRING) WITH ("
+                                + "'connector' = 'doris',"
+                                + "'fenodes' = '%s',"
+                                + "'jdbc-url' = '%s',"
+                                + "'table.identifier' = '%s.%s',"
+                                + "'username' = '%s',"
+                                + "'password' = '%s',"
+                                + "'sink.write-mode' = 'TVF',"
+                                + "'sink.parallelism' = '1',"
+                                + "'sink.label-prefix' = '%s',"
+                                + "'sink.s3.endpoint' = '%s',"
+                                + "'sink.s3.region' = '%s',"
+                                + "'sink.s3.bucket' = '%s',"
+                                + "'sink.s3.prefix' = '%s',"
+                                + "'sink.s3.access-key' = '%s',"
+                                + "'sink.s3.secret-key' = '%s',"
+                                + "'sink.s3.path-style-access' = 'true',"
+                                + "'sink.properties.compress_type' = '')",
+                        getFenodes(),
+                        getDorisQueryUrl(),
+                        DATABASE,
+                        table,
+                        getDorisUsername(),
+                        getDorisPassword(),
+                        labelPrefix,
+                        s3Endpoint,
+                        REGION,
+                        BUCKET,
+                        objectPrefix,
+                        ACCESS_KEY,
+                        SECRET_KEY));
+        TableResult result = tableEnv.executeSql("INSERT INTO tvf_sink VALUES 
(1, 'doris')");
+        waitForJobStatus(
+                result.getJobClient().get(),
+                Collections.singletonList(FINISHED),
+                Deadline.fromNow(Duration.ofSeconds(120)));
+
+        assertResult(table, "id,name", Collections.singletonList("1,doris"), 
2);
+        ListObjectsV2Response objects =
+                s3Client.listObjectsV2(
+                        ListObjectsV2Request.builder()
+                                .bucket(BUCKET)
+                                .prefix(objectPrefix + "/" + labelPrefix + "_" 
+ table)
+                                .build());
+        Assert.assertTrue("Expected an uncompressed TVF object", 
objects.keyCount() > 0);
+        Assert.assertTrue(
+                "Expected TVF objects to remain uncompressed",
+                objects.contents().stream().allMatch(object -> 
object.key().endsWith(".json")));
     }
 
     @Test


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

Reply via email to