This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new b9eacfae2 [INLONG-1517][Manager] Support sink data to ClickHouse 
(#3625)
b9eacfae2 is described below

commit b9eacfae2f08484a047e87f204cfa0b0c7ee69ef
Author: healchow <[email protected]>
AuthorDate: Mon Apr 11 19:29:39 2022 +0800

    [INLONG-1517][Manager] Support sink data to ClickHouse (#3625)
---
 .../manager/client/api/sink/ClickHouseSink.java    | 31 ++++++------
 .../client/api/util/InlongStreamSinkTransfer.java  | 24 +++++-----
 .../common/pojo/sink/ck/ClickHouseSinkDTO.java     | 55 +++++++++++-----------
 .../pojo/sink/ck/ClickHouseSinkListResponse.java   | 39 +++++++--------
 .../common/pojo/sink/ck/ClickHouseSinkRequest.java | 43 ++++++++---------
 .../pojo/sink/ck/ClickHouseSinkResponse.java       | 48 ++++++++++++-------
 .../thirdparty/sort/util/SinkInfoUtils.java        | 37 ++++++++-------
 .../core/sink/ClickHouseStreamSinkServiceTest.java |  2 +-
 8 files changed, 146 insertions(+), 133 deletions(-)

diff --git 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/sink/ClickHouseSink.java
 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/sink/ClickHouseSink.java
index b507b9e59..2d76714e3 100644
--- 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/sink/ClickHouseSink.java
+++ 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/sink/ClickHouseSink.java
@@ -45,7 +45,7 @@ public class ClickHouseSink extends StreamSink {
     private String jdbcUrl;
 
     @ApiModelProperty("Target database name")
-    private String databaseName;
+    private String dbName;
 
     @ApiModelProperty("Target table name")
     private String tableName;
@@ -53,26 +53,27 @@ public class ClickHouseSink extends StreamSink {
     @ApiModelProperty("Authentication for clickhouse")
     private DefaultAuthentication authentication;
 
-    @ApiModelProperty("Whether distributed table")
-    private Boolean distributedTable;
+    @ApiModelProperty("Flush interval, unit: second, default is 1s")
+    private Integer flushInterval;
 
-    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
-    private String partitionStrategy;
+    @ApiModelProperty("Flush when record number reaches flushRecord")
+    private Integer flushRecord;
 
-    @ApiModelProperty("Partition key")
-    private String partitionKey;
+    @ApiModelProperty("Write max retry times, default is 3")
+    private Integer retryTimes;
 
-    @ApiModelProperty("Key field names")
-    private String[] keyFieldNames;
+    @ApiModelProperty("Whether distributed table? 0: no, 1: yes")
+    private Integer isDistributed;
 
-    @ApiModelProperty("Flush interval")
-    private Integer flushInterval;
+    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
+    private String partitionStrategy;
 
-    @ApiModelProperty("Flush record number")
-    private Integer flushRecordNumber;
+    @ApiModelProperty(value = "Partition files, separate with commas",
+            notes = "Necessary when partitionStrategy is HASH, must be one of 
the field list")
+    private String partitionFields;
 
-    @ApiModelProperty("Write max retry times")
-    private Integer writeMaxRetryTimes;
+    @ApiModelProperty("Key field names, separate with commas")
+    private String keyFieldNames;
 
     @ApiModelProperty("Create topic or not")
     private boolean needCreated;
diff --git 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/util/InlongStreamSinkTransfer.java
 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/util/InlongStreamSinkTransfer.java
index 07c7bd051..590b9fbae 100644
--- 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/util/InlongStreamSinkTransfer.java
+++ 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/util/InlongStreamSinkTransfer.java
@@ -89,22 +89,22 @@ public class InlongStreamSinkTransfer {
         ClickHouseSinkRequest clickHouseSinkRequest = new 
ClickHouseSinkRequest();
         ClickHouseSink clickHouseSink = (ClickHouseSink) streamSink;
         clickHouseSinkRequest.setSinkName(clickHouseSink.getSinkName());
-        
clickHouseSinkRequest.setDatabaseName(clickHouseSink.getDatabaseName());
+        clickHouseSinkRequest.setDbName(clickHouseSink.getDbName());
         clickHouseSinkRequest.setSinkType(clickHouseSink.getSinkType().name());
         clickHouseSinkRequest.setJdbcUrl(clickHouseSink.getJdbcUrl());
         DefaultAuthentication defaultAuthentication = 
clickHouseSink.getAuthentication();
         AssertUtil.notNull(defaultAuthentication,
-                String.format("Clickhouse storage:%s must be authenticated", 
clickHouseSink.getDatabaseName()));
+                String.format("Clickhouse storage:%s must be authenticated", 
clickHouseSink.getDbName()));
         clickHouseSinkRequest.setUsername(defaultAuthentication.getUserName());
         clickHouseSinkRequest.setPassword(defaultAuthentication.getPassword());
         clickHouseSinkRequest.setTableName(clickHouseSink.getTableName());
-        
clickHouseSinkRequest.setDistributedTable(clickHouseSink.getDistributedTable());
+        
clickHouseSinkRequest.setIsDistributed(clickHouseSink.getIsDistributed());
         
clickHouseSinkRequest.setFlushInterval(clickHouseSink.getFlushInterval());
-        
clickHouseSinkRequest.setFlushRecordNumber(clickHouseSink.getFlushRecordNumber());
+        clickHouseSinkRequest.setFlushRecord(clickHouseSink.getFlushRecord());
         
clickHouseSinkRequest.setKeyFieldNames(clickHouseSink.getKeyFieldNames());
-        
clickHouseSinkRequest.setPartitionKey(clickHouseSink.getPartitionKey());
+        
clickHouseSinkRequest.setPartitionFields(clickHouseSink.getPartitionFields());
         
clickHouseSinkRequest.setPartitionStrategy(clickHouseSink.getPartitionStrategy());
-        
clickHouseSinkRequest.setWriteMaxRetryTimes(clickHouseSink.getWriteMaxRetryTimes());
+        clickHouseSinkRequest.setRetryTimes(clickHouseSink.getRetryTimes());
         clickHouseSinkRequest.setInlongGroupId(streamInfo.getInlongGroupId());
         
clickHouseSinkRequest.setInlongStreamId(streamInfo.getInlongStreamId());
         clickHouseSinkRequest.setProperties(clickHouseSink.getProperties());
@@ -125,19 +125,19 @@ public class InlongStreamSinkTransfer {
             ClickHouseSink snapshot = (ClickHouseSink) streamSink;
             clickHouseSink = CommonBeanUtils.copyProperties(snapshot, 
ClickHouseSink::new);
         } else {
-            
clickHouseSink.setDistributedTable(sinkResponse.getDistributedTable());
+            clickHouseSink.setIsDistributed(sinkResponse.getIsDistributed());
             clickHouseSink.setSinkName(sinkResponse.getSinkName());
             clickHouseSink.setFlushInterval(sinkResponse.getFlushInterval());
             clickHouseSink.setAuthentication(new 
DefaultAuthentication(sinkResponse.getSinkName(),
                     sinkResponse.getPassword()));
-            clickHouseSink.setDatabaseName(sinkResponse.getDatabaseName());
-            
clickHouseSink.setFlushRecordNumber(sinkResponse.getFlushRecordNumber());
+            clickHouseSink.setDbName(sinkResponse.getDbName());
+            clickHouseSink.setFlushRecord(sinkResponse.getFlushRecord());
             clickHouseSink.setJdbcUrl(sinkResponse.getJdbcUrl());
-            clickHouseSink.setPartitionKey(sinkResponse.getPartitionKey());
+            
clickHouseSink.setPartitionFields(sinkResponse.getPartitionFields());
             clickHouseSink.setKeyFieldNames(sinkResponse.getKeyFieldNames());
             
clickHouseSink.setPartitionStrategy(sinkResponse.getPartitionStrategy());
-            
clickHouseSink.setWriteMaxRetryTimes(sinkResponse.getWriteMaxRetryTimes());
-            
clickHouseSink.setDistributedTable(sinkResponse.getDistributedTable());
+            clickHouseSink.setRetryTimes(sinkResponse.getRetryTimes());
+            clickHouseSink.setIsDistributed(sinkResponse.getIsDistributed());
         }
         clickHouseSink.setProperties(sinkResponse.getProperties());
         clickHouseSink.setNeedCreated(sinkResponse.getEnableCreateResource() 
== 1);
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkDTO.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkDTO.java
index bce8e709e..2efecea8a 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkDTO.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkDTO.java
@@ -20,7 +20,6 @@ package org.apache.inlong.manager.common.pojo.sink.ck;
 import com.fasterxml.jackson.databind.DeserializationFeature;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import io.swagger.annotations.ApiModelProperty;
-import java.util.Map;
 import lombok.AllArgsConstructor;
 import lombok.Builder;
 import lombok.Data;
@@ -29,6 +28,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 
 import javax.validation.constraints.NotNull;
+import java.util.Map;
 
 /**
  * Sink info of ClickHouse
@@ -44,38 +44,39 @@ public class ClickHouseSinkDTO {
     @ApiModelProperty("ClickHouse JDBC URL")
     private String jdbcUrl;
 
-    @ApiModelProperty("Target database name")
-    private String databaseName;
-
-    @ApiModelProperty("Target table name")
-    private String tableName;
-
     @ApiModelProperty("Username for JDBC URL")
     private String username;
 
     @ApiModelProperty("User password")
     private String password;
 
-    @ApiModelProperty("Whether distributed table")
-    private Boolean distributedTable;
+    @ApiModelProperty("Target database name")
+    private String dbName;
 
-    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
-    private String partitionStrategy;
+    @ApiModelProperty("Target table name")
+    private String tableName;
 
-    @ApiModelProperty("Partition key")
-    private String partitionKey;
+    @ApiModelProperty("Flush interval, unit: second, default is 1s")
+    private Integer flushInterval;
 
-    @ApiModelProperty("Key field names")
-    private String[] keyFieldNames;
+    @ApiModelProperty("Flush when record number reaches flushRecord")
+    private Integer flushRecord;
 
-    @ApiModelProperty("Flush interval")
-    private Integer flushInterval;
+    @ApiModelProperty("Write max retry times, default is 3")
+    private Integer retryTimes;
 
-    @ApiModelProperty("Flush record number")
-    private Integer flushRecordNumber;
+    @ApiModelProperty("Whether distributed table? 0: no, 1: yes")
+    private Integer isDistributed;
 
-    @ApiModelProperty("Write max retry times")
-    private Integer writeMaxRetryTimes;
+    @ApiModelProperty("Partition strategy, support: BALANCE, RANDOM, HASH")
+    private String partitionStrategy;
+
+    @ApiModelProperty(value = "Partition files, separate with commas",
+            notes = "Necessary when partitionStrategy is HASH, must be one of 
the field list")
+    private String partitionFields;
+
+    @ApiModelProperty("Key field names, separate with commas")
+    private String keyFieldNames;
 
     @ApiModelProperty("Properties for clickhouse")
     private Map<String, Object> properties;
@@ -88,15 +89,15 @@ public class ClickHouseSinkDTO {
                 .jdbcUrl(request.getJdbcUrl())
                 .username(request.getUsername())
                 .password(request.getPassword())
-                .databaseName(request.getDatabaseName())
+                .dbName(request.getDbName())
                 .tableName(request.getTableName())
-                .distributedTable(request.getDistributedTable())
+                .flushInterval(request.getFlushInterval())
+                .flushRecord(request.getFlushRecord())
+                .retryTimes(request.getRetryTimes())
+                .isDistributed(request.getIsDistributed())
                 .partitionStrategy(request.getPartitionStrategy())
-                .partitionKey(request.getPartitionKey())
+                .partitionFields(request.getPartitionFields())
                 .keyFieldNames(request.getKeyFieldNames())
-                .flushInterval(request.getFlushInterval())
-                .flushRecordNumber(request.getFlushRecordNumber())
-                .writeMaxRetryTimes(request.getWriteMaxRetryTimes())
                 .properties(request.getProperties())
                 .build();
     }
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkListResponse.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkListResponse.java
index 63b1db52e..f4cb3bf5d 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkListResponse.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkListResponse.java
@@ -34,37 +34,32 @@ public class ClickHouseSinkListResponse extends 
SinkListResponse {
     @ApiModelProperty("ClickHouse JDBC URL")
     private String jdbcUrl;
 
-    @ApiModelProperty("Target database name")
-    private String databaseName;
-
-    @ApiModelProperty("Target table name")
-    private String tableName;
-
     @ApiModelProperty("Username for JDBC URL")
     private String username;
 
-    @ApiModelProperty("User password")
-    private String password;
+    @ApiModelProperty("Target database name")
+    private String dbName;
 
-    @ApiModelProperty("Whether distributed table")
-    private Boolean distributedTable;
+    @ApiModelProperty("Target table name")
+    private String tableName;
 
-    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
-    private String partitionStrategy;
+    @ApiModelProperty("Flush interval, unit: second, default is 1s")
+    private Integer flushInterval;
 
-    @ApiModelProperty("Partition key")
-    private String partitionKey;
+    @ApiModelProperty("Flush when record number reaches flushRecord")
+    private Integer flushRecord;
 
-    @ApiModelProperty("Key field names")
-    private String[] keyFieldNames;
+    @ApiModelProperty("Write max retry times, default is 3")
+    private Integer retryTimes;
 
-    @ApiModelProperty("Flush interval")
-    private Integer flushInterval;
+    @ApiModelProperty("Whether distributed table? 0: no, 1: yes")
+    private Integer isDistributed;
 
-    @ApiModelProperty("Flush record number")
-    private Integer flushRecordNumber;
+    @ApiModelProperty("Partition strategy, support: BALANCE, RANDOM, HASH")
+    private String partitionStrategy;
 
-    @ApiModelProperty("Write max retry times")
-    private Integer writeMaxRetryTimes;
+    @ApiModelProperty(value = "Partition files, separate with commas",
+            notes = "Necessary when partitionStrategy is HASH, must be one of 
the field list")
+    private String partitionFields;
 
 }
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkRequest.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkRequest.java
index 6493ae5c3..0d35fcf91 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkRequest.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkRequest.java
@@ -27,7 +27,7 @@ import org.apache.inlong.manager.common.pojo.sink.SinkRequest;
 import org.apache.inlong.manager.common.util.JsonTypeDefine;
 
 /**
- * Request of the ClickHouse sink info
+ * Request of the ClickHouse sink.
  */
 @Data
 @ToString(callSuper = true)
@@ -39,37 +39,38 @@ public class ClickHouseSinkRequest extends SinkRequest {
     @ApiModelProperty("ClickHouse JDBC URL")
     private String jdbcUrl;
 
-    @ApiModelProperty("Target database name")
-    private String databaseName;
-
-    @ApiModelProperty("Target table name")
-    private String tableName;
-
     @ApiModelProperty("Username for JDBC URL")
     private String username;
 
     @ApiModelProperty("User password")
     private String password;
 
-    @ApiModelProperty("Whether distributed table")
-    private Boolean distributedTable;
+    @ApiModelProperty("Target database name")
+    private String dbName;
 
-    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
-    private String partitionStrategy;
+    @ApiModelProperty("Target table name")
+    private String tableName;
 
-    @ApiModelProperty("Partition key")
-    private String partitionKey;
+    @ApiModelProperty("Flush interval, unit: second, default is 1s")
+    private Integer flushInterval;
 
-    @ApiModelProperty("Key field names")
-    private String[] keyFieldNames;
+    @ApiModelProperty("Flush when record number reaches flushRecord")
+    private Integer flushRecord;
 
-    @ApiModelProperty("Flush interval")
-    private Integer flushInterval;
+    @ApiModelProperty("Write max retry times, default is 3")
+    private Integer retryTimes;
+
+    @ApiModelProperty("Whether distributed table? 0: no, 1: yes")
+    private Integer isDistributed;
+
+    @ApiModelProperty("Partition strategy, support: BALANCE, RANDOM, HASH")
+    private String partitionStrategy;
 
-    @ApiModelProperty("Flush record number")
-    private Integer flushRecordNumber;
+    @ApiModelProperty(value = "Partition files, separate with commas",
+            notes = "Necessary when partitionStrategy is HASH, must be one of 
the field list")
+    private String partitionFields;
 
-    @ApiModelProperty("Write max retry times")
-    private Integer writeMaxRetryTimes;
+    @ApiModelProperty("Key field names, separate with commas")
+    private String keyFieldNames;
 
 }
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkResponse.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkResponse.java
index bb0e47654..21659253b 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkResponse.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/ck/ClickHouseSinkResponse.java
@@ -26,7 +26,7 @@ import org.apache.inlong.manager.common.enums.SinkType;
 import org.apache.inlong.manager.common.pojo.sink.SinkResponse;
 
 /**
- * Response of the ClickHouse sink
+ * Response of the ClickHouse sink.
  */
 @Data
 @ToString(callSuper = true)
@@ -36,28 +36,40 @@ public class ClickHouseSinkResponse extends SinkResponse {
 
     @ApiModelProperty("ClickHouse JDBC URL")
     private String jdbcUrl;
-    @ApiModelProperty("Target database name")
-    private String databaseName;
-    @ApiModelProperty("Target table name")
-    private String tableName;
+
     @ApiModelProperty("Username for JDBC URL")
     private String username;
+
     @ApiModelProperty("User password")
     private String password;
-    @ApiModelProperty("Whether distributed table")
-    private Boolean distributedTable;
-    @ApiModelProperty("Partition strategy,support: BALANCE, RANDOM, HASH")
-    private String partitionStrategy;
-    @ApiModelProperty("Partition key")
-    private String partitionKey;
-    @ApiModelProperty("Key field names")
-    private String[] keyFieldNames;
-    @ApiModelProperty("Flush interval")
+
+    @ApiModelProperty("Target database name")
+    private String dbName;
+
+    @ApiModelProperty("Target table name")
+    private String tableName;
+
+    @ApiModelProperty("Flush interval, unit: second, default is 1s")
     private Integer flushInterval;
-    @ApiModelProperty("Flush record number")
-    private Integer flushRecordNumber;
-    @ApiModelProperty("Write max retry times")
-    private Integer writeMaxRetryTimes;
+
+    @ApiModelProperty("Flush when record number reaches flushRecord")
+    private Integer flushRecord;
+
+    @ApiModelProperty("Write max retry times, default is 3")
+    private Integer retryTimes;
+
+    @ApiModelProperty("Whether distributed table? 0: no, 1: yes")
+    private Integer isDistributed;
+
+    @ApiModelProperty("Partition strategy, support: BALANCE, RANDOM, HASH")
+    private String partitionStrategy;
+
+    @ApiModelProperty(value = "Partition files, separate with commas",
+            notes = "Necessary when partitionStrategy is HASH, must be one of 
the field list")
+    private String partitionFields;
+
+    @ApiModelProperty("Key field names, separate with commas")
+    private String keyFieldNames;
 
     public ClickHouseSinkResponse() {
         this.sinkType = SinkType.SINK_CLICKHOUSE;
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/thirdparty/sort/util/SinkInfoUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/thirdparty/sort/util/SinkInfoUtils.java
index d1490c020..0a06f59f9 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/thirdparty/sort/util/SinkInfoUtils.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/thirdparty/sort/util/SinkInfoUtils.java
@@ -86,30 +86,33 @@ public class SinkInfoUtils {
             throw new BusinessException(String.format("ClickHouse={%s} fields 
cannot be empty", sinkResponse));
         } else if (StringUtils.isEmpty(sinkResponse.getTableName())) {
             throw new BusinessException(String.format("ClickHouse={%s} table 
name cannot be empty", sinkResponse));
-        } else if (StringUtils.isEmpty(sinkResponse.getDatabaseName())) {
+        } else if (StringUtils.isEmpty(sinkResponse.getDbName())) {
             throw new BusinessException(String.format("ClickHouse={%s} 
database name cannot be empty", sinkResponse));
         }
-        if (sinkResponse.getDistributedTable() == null) {
-            throw new BusinessException(String.format("ClickHouse={%s} 
distribute cannot be empty", sinkResponse));
+
+        Integer isDistributed = sinkResponse.getIsDistributed();
+        if (isDistributed == null) {
+            throw new BusinessException(String.format("ClickHouse={%s} 
isDistributed cannot be null", sinkResponse));
         }
 
-        ClickHouseSinkInfo.PartitionStrategy partitionStrategy;
-        if 
(PartitionStrategy.BALANCE.name().equalsIgnoreCase(sinkResponse.getPartitionStrategy()))
 {
-            partitionStrategy = PartitionStrategy.BALANCE;
-        } else if 
(PartitionStrategy.HASH.name().equalsIgnoreCase(sinkResponse.getPartitionStrategy()))
 {
-            partitionStrategy = PartitionStrategy.HASH;
-        } else if 
(PartitionStrategy.RANDOM.name().equalsIgnoreCase(sinkResponse.getPartitionStrategy()))
 {
-            partitionStrategy = PartitionStrategy.RANDOM;
-        } else {
-            partitionStrategy = PartitionStrategy.RANDOM;
+        // Default partition strategy is RANDOM
+        ClickHouseSinkInfo.PartitionStrategy partitionStrategy = 
PartitionStrategy.RANDOM;
+        boolean distributedTable = isDistributed == 1;
+        if (distributedTable) {
+            if 
(PartitionStrategy.BALANCE.name().equalsIgnoreCase(sinkResponse.getPartitionStrategy()))
 {
+                partitionStrategy = PartitionStrategy.BALANCE;
+            } else if 
(PartitionStrategy.HASH.name().equalsIgnoreCase(sinkResponse.getPartitionStrategy()))
 {
+                partitionStrategy = PartitionStrategy.HASH;
+            }
         }
 
-        return new ClickHouseSinkInfo(sinkResponse.getJdbcUrl(), 
sinkResponse.getDatabaseName(),
+        // TODO Add keyFieldNames instead of `new String[0]`
+        return new ClickHouseSinkInfo(sinkResponse.getJdbcUrl(), 
sinkResponse.getDbName(),
                 sinkResponse.getTableName(), sinkResponse.getUsername(), 
sinkResponse.getPassword(),
-                sinkResponse.getDistributedTable(), partitionStrategy, 
sinkResponse.getPartitionKey(),
-                sinkFields.toArray(new FieldInfo[0]), 
sinkResponse.getKeyFieldNames(),
-                sinkResponse.getFlushInterval(), 
sinkResponse.getFlushRecordNumber(),
-                sinkResponse.getWriteMaxRetryTimes());
+                distributedTable, partitionStrategy, 
sinkResponse.getPartitionFields(),
+                sinkFields.toArray(new FieldInfo[0]), new String[0],
+                sinkResponse.getFlushInterval(), sinkResponse.getFlushRecord(),
+                sinkResponse.getRetryTimes());
     }
 
     // TODO Need set more configs for IcebergSinkInfo
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/sink/ClickHouseStreamSinkServiceTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/sink/ClickHouseStreamSinkServiceTest.java
index 0484cd74f..98cd24420 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/sink/ClickHouseStreamSinkServiceTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/sink/ClickHouseStreamSinkServiceTest.java
@@ -62,7 +62,7 @@ public class ClickHouseStreamSinkServiceTest extends 
ServiceBaseTest {
         sinkInfo.setSinkType(SinkType.SINK_CLICKHOUSE);
         sinkInfo.setJdbcUrl(ckJdbcUrl);
         sinkInfo.setUsername(ckUsername);
-        sinkInfo.setDatabaseName(ckDatabaseName);
+        sinkInfo.setDbName(ckDatabaseName);
         sinkInfo.setTableName(ckTableName);
         sinkInfo.setEnableCreateResource(Constant.DISABLE_CREATE_RESOURCE);
         sinkId = sinkService.save(sinkInfo, globalOperator);

Reply via email to