1
0

[HUDI-2891] Fix write configs for Java engine in Kafka Connect Sink (#4161)

This commit is contained in:
Y Ethan Guo
2021-11-30 06:45:50 -08:00
committed by GitHub
parent a398aad1fc
commit ea009b55a3
4 changed files with 4 additions and 2 deletions

View File

@@ -10,7 +10,6 @@
"topics": "hudi-test-topic",
"hoodie.table.name": "hudi-test-topic",
"hoodie.table.type": "MERGE_ON_READ",
"hoodie.metadata.enable": "false",
"hoodie.base.path": "hdfs://namenode:8020/user/hive/warehouse/hudi-test-topic",
"hoodie.datasource.write.recordkey.field": "volume",
"hoodie.datasource.write.partitionpath.field": "date",

View File

@@ -10,7 +10,6 @@
"topics": "hudi-test-topic",
"hoodie.table.name": "hudi-test-topic",
"hoodie.table.type": "MERGE_ON_READ",
"hoodie.metadata.enable": "false",
"hoodie.base.path": "file:///tmp/hoodie/hudi-test-topic",
"hoodie.datasource.write.recordkey.field": "volume",
"hoodie.datasource.write.partitionpath.field": "date",

View File

@@ -22,6 +22,7 @@ import org.apache.hudi.client.HoodieJavaWriteClient;
import org.apache.hudi.client.WriteStatus;
import org.apache.hudi.client.common.HoodieJavaEngineContext;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.EngineType;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.model.HoodieAvroPayload;
import org.apache.hudi.common.util.ReflectionUtils;
@@ -74,6 +75,7 @@ public class KafkaConnectWriterProvider implements ConnectWriterProvider<WriteSt
// Create the write client to write some records in
writeConfig = HoodieWriteConfig.newBuilder()
.withEngineType(EngineType.JAVA)
.withProperties(connectConfigs.getProps())
.withFileIdPrefixProviderClassName(KafkaConnectFileIdPrefixProvider.class.getName())
.withProps(Collections.singletonMap(

View File

@@ -20,6 +20,7 @@ package org.apache.hudi.writers;
import org.apache.hudi.client.HoodieJavaWriteClient;
import org.apache.hudi.client.common.HoodieJavaEngineContext;
import org.apache.hudi.common.engine.EngineType;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
import org.apache.hudi.common.util.Option;
@@ -62,6 +63,7 @@ public class TestBufferedConnectWriter {
configs = KafkaConnectConfigs.newBuilder().build();
schemaProvider = new TestAbstractConnectWriter.TestSchemaProvider();
writeConfig = HoodieWriteConfig.newBuilder()
.withEngineType(EngineType.JAVA)
.withPath("/tmp")
.withSchema(schemaProvider.getSourceSchema().toString())
.build();