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]