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/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 65cb315b5 [INLONG-7711][Sort][Manager] Support specifying parameters 
for the kudu client (#7718)
65cb315b5 is described below

commit 65cb315b5bc432bfa5eb45431e1f29c12740295d
Author: feat <[email protected]>
AuthorDate: Thu Mar 30 15:43:30 2023 +0800

    [INLONG-7711][Sort][Manager] Support specifying parameters for the kudu 
client (#7718)
---
 .../manager/pojo/node/kudu/KuduDataNodeDTO.java    | 12 +++++-----
 .../manager/pojo/node/kudu/KuduDataNodeInfo.java   | 12 +++++-----
 .../pojo/node/kudu/KuduDataNodeRequest.java        |  6 ++---
 .../inlong/sort/kudu/common/KuduOptions.java       | 27 ++++++++++++++++++++++
 .../sort/kudu/sink/AbstractKuduSinkFunction.java   | 16 +++++++++++++
 .../sort/kudu/sink/KuduAsyncSinkFunction.java      |  3 +--
 .../inlong/sort/kudu/sink/KuduSinkFunction.java    |  3 +--
 .../sort/kudu/source/KuduLookupFunction.java       | 19 ++++++++++++++-
 .../sort/kudu/table/KuduDynamicTableFactory.java   |  8 +++++++
 9 files changed, 86 insertions(+), 20 deletions(-)

diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeDTO.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeDTO.java
index 4e9a59300..fcbdc0315 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeDTO.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeDTO.java
@@ -47,17 +47,17 @@ public class KuduDataNodeDTO {
     @ApiModelProperty("Kudu masters, a comma separated list of 'host:port' 
pairs")
     private String masters;
 
-    @ApiModelProperty("Default admin operation timeout in ms, default is 3000")
+    @ApiModelProperty("Sets the default timeout used for administrative 
operations (e.g. createTable, deleteTable, etc). Optional. If not provided, 
defaults to 30s. A value of 0 disables the timeout")
     private Integer defaultAdminOperationTimeoutMs;
 
-    @ApiModelProperty("Default operation timeout in ms, default is 3000")
-    private Integer defaultOperationTimeoutMs = 3000;
+    @ApiModelProperty("Sets the default timeout used for user operations 
(using sessions and scanners). Optional. If not provided, defaults to 30s. A 
value of 0 disables the timeout")
+    private Integer defaultOperationTimeoutMs = 30000;
 
     @ApiModelProperty("Default socket read timeout in ms, default is 10000")
-    private Integer defaultSocketReadTimeoutMs;
+    private Integer defaultSocketReadTimeoutMs = 10000;
 
-    @ApiModelProperty("Whether to enable the statistics collection function of 
the Kudu client, default is false.")
-    private Boolean statisticsDisabled;
+    @ApiModelProperty("Disable this client's collection of statistics. 
Statistics are enabled by default")
+    private Boolean statisticsDisabled = false;
 
     /**
      * Get the dto instance from the request
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeInfo.java
index 307123809..a73625c60 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeInfo.java
@@ -46,17 +46,17 @@ public class KuduDataNodeInfo extends DataNodeInfo {
     @ApiModelProperty("Kudu masters, a comma separated list of 'host:port' 
pairs")
     private String masters;
 
-    @ApiModelProperty("Default admin operation timeout in ms, default is 3000")
+    @ApiModelProperty("Sets the default timeout used for administrative 
operations (e.g. createTable, deleteTable, etc). Optional. If not provided, 
defaults to 30s. A value of 0 disables the timeout")
     private Integer defaultAdminOperationTimeoutMs;
 
-    @ApiModelProperty("Default operation timeout in ms, default is 3000")
-    private Integer defaultOperationTimeoutMs = 3000;
+    @ApiModelProperty("Sets the default timeout used for user operations 
(using sessions and scanners). Optional. If not provided, defaults to 30s. A 
value of 0 disables the timeout")
+    private Integer defaultOperationTimeoutMs = 30000;
 
     @ApiModelProperty("Default socket read timeout in ms, default is 10000")
-    private Integer defaultSocketReadTimeoutMs;
+    private Integer defaultSocketReadTimeoutMs = 10000;
 
-    @ApiModelProperty("Whether to enable the statistics collection function of 
the Kudu client, default is false.")
-    private Boolean statisticsDisabled;
+    @ApiModelProperty("Disable this client's collection of statistics. 
Statistics are enabled by default")
+    private Boolean statisticsDisabled = false;
 
     public KuduDataNodeInfo() {
         setType(DataNodeType.KUDU);
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeRequest.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeRequest.java
index 4893df222..794a263aa 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeRequest.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/node/kudu/KuduDataNodeRequest.java
@@ -40,16 +40,16 @@ public class KuduDataNodeRequest extends DataNodeRequest {
     @ApiModelProperty("Kudu masters, a comma separated list of 'host:port' 
pairs")
     private String masters;
 
-    @ApiModelProperty("Default admin operation timeout in ms, default is 
30000")
+    @ApiModelProperty("Sets the default timeout used for user operations 
(using sessions and scanners). Optional. If not provided, defaults to 30s. A 
value of 0 disables the timeout")
     private Integer defaultAdminOperationTimeoutMs = 30000;
 
-    @ApiModelProperty("Default operation timeout in ms, default is 30000")
+    @ApiModelProperty("Sets the default timeout used for user operations 
(using sessions and scanners). Optional. If not provided, defaults to 30s. A 
value of 0 disables the timeout")
     private Integer defaultOperationTimeoutMs = 30000;
 
     @ApiModelProperty("Default socket read timeout in ms, default is 10000")
     private Integer defaultSocketReadTimeoutMs = 10000;
 
-    @ApiModelProperty("Whether to enable the statistics collection function of 
the Kudu client, default is false.")
+    @ApiModelProperty("Disable this client's collection of statistics. 
Statistics are enabled by default")
     private Boolean statisticsDisabled = false;
 
     public KuduDataNodeRequest() {
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/common/KuduOptions.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/common/KuduOptions.java
index 2a964e4ef..4ef0f4472 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/common/KuduOptions.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/common/KuduOptions.java
@@ -49,6 +49,32 @@ public class KuduOptions {
                     .withDescription("The maximum number of results cached in 
the " +
                             "lookup source.");
 
+    public static final ConfigOption<Long> 
DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS =
+            ConfigOptions.key("default-admin-operation-timeout")
+                    .longType()
+                    .defaultValue(30000L)
+                    .withDescription(
+                            "Sets the default timeout used for administrative 
operations (e.g. createTable, deleteTable, etc). Optional. If not provided, 
defaults to 30s. A value of 0 disables the timeout.");
+
+    public static final ConfigOption<Long> DEFAULT_OPERATION_TIMEOUT_IN_MS =
+            ConfigOptions.key("default-operation-timeout")
+                    .longType()
+                    .defaultValue(30000L)
+                    .withDescription(
+                            "Sets the default timeout used for user operations 
(using sessions and scanners). Optional. If not provided, defaults to 30s. A 
value of 0 disables the timeout.");
+
+    public static final ConfigOption<Long> DEFAULT_SOCKET_READ_TIMEOUT_IN_MS =
+            ConfigOptions.key("default-socket-read-timeout")
+                    .longType()
+                    .defaultValue(10000L)
+                    .withDescription("Default socket read timeout in ms, 
default is 10000");
+    public static final ConfigOption<Boolean> DISABLED_STATISTICS =
+            ConfigOptions.key("disabled-statistics")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "Disable this client's collection of statistics. 
Statistics are enabled by default.");
+
     public static final ConfigOption<String> MAX_CACHE_TIME =
             ConfigOptions.key("lookup.max-cache-time")
                     .stringType()
@@ -67,6 +93,7 @@ public class KuduOptions {
                     .defaultValue(3)
                     .withDescription("The maximum number of retries when an " +
                             "exception is caught.");
+
     public static final ConfigOption<Integer> MAX_BUFFER_SIZE =
             ConfigOptions.key("sink.max-buffer-size")
                     .intType()
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/AbstractKuduSinkFunction.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/AbstractKuduSinkFunction.java
index 6acb9cd32..d092b7370 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/AbstractKuduSinkFunction.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/AbstractKuduSinkFunction.java
@@ -33,6 +33,7 @@ import org.apache.inlong.sort.base.metric.MetricState;
 import org.apache.inlong.sort.base.metric.SinkMetricData;
 import org.apache.inlong.sort.base.util.MetricStateUtils;
 import org.apache.inlong.sort.kudu.common.KuduTableInfo;
+import org.apache.kudu.client.KuduClient;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -43,8 +44,12 @@ import static 
org.apache.inlong.sort.base.Constants.DIRTY_RECORDS_OUT;
 import static org.apache.inlong.sort.base.Constants.INLONG_METRIC_STATE_NAME;
 import static org.apache.inlong.sort.base.Constants.NUM_BYTES_OUT;
 import static org.apache.inlong.sort.base.Constants.NUM_RECORDS_OUT;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_OPERATION_TIMEOUT_IN_MS;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_SOCKET_READ_TIMEOUT_IN_MS;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_BUFFER_SIZE;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_RETRIES;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DISABLED_STATISTICS;
 
 /**
  * The base for all kudu sinks.
@@ -126,7 +131,18 @@ public abstract class AbstractKuduSinkFunction
             sinkMetricData = new SinkMetricData(metricOption, 
getRuntimeContext().getMetricGroup());
         }
     }
+    protected KuduClient buildKuduClient() {
+        KuduClient.KuduClientBuilder builder = new 
KuduClient.KuduClientBuilder(kuduTableInfo.getMasters());
+        if (configuration.getBoolean(DISABLED_STATISTICS)) {
+            builder.disableStatistics();
+        }
+        
builder.defaultAdminOperationTimeoutMs(configuration.getLong(DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS));
+        
builder.defaultOperationTimeoutMs(configuration.getLong(DEFAULT_OPERATION_TIMEOUT_IN_MS));
+        
builder.defaultSocketReadTimeoutMs(configuration.getLong(DEFAULT_SOCKET_READ_TIMEOUT_IN_MS));
 
+        return builder
+                .build();
+    }
     @Override
     public void invoke(RowData row, Context context) throws Exception {
         addBatch(row);
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduAsyncSinkFunction.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduAsyncSinkFunction.java
index 47eda9f1b..b816e1e7a 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduAsyncSinkFunction.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduAsyncSinkFunction.java
@@ -99,8 +99,7 @@ public class KuduAsyncSinkFunction
                 namedThreadFactory,
                 new ThreadPoolExecutor.AbortPolicy());
 
-        KuduClient client = new 
KuduClient.KuduClientBuilder(kuduTableInfo.getMasters())
-                .build();
+        KuduClient client = buildKuduClient();
         KuduTable kuduTable;
         try {
             String tableName = kuduTableInfo.getTableName();
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduSinkFunction.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduSinkFunction.java
index c06f05643..63233d033 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduSinkFunction.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/sink/KuduSinkFunction.java
@@ -66,8 +66,7 @@ public class KuduSinkFunction
 
         boolean forceWithUpsertMode = 
configuration.getBoolean(SINK_FORCE_WITH_UPSERT_MODE);
 
-        KuduClient client = new 
KuduClient.KuduClientBuilder(kuduTableInfo.getMasters())
-                .build();
+        KuduClient client = buildKuduClient();
         KuduTable kuduTable;
         try {
             String tableName = kuduTableInfo.getTableName();
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/source/KuduLookupFunction.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/source/KuduLookupFunction.java
index 626099950..ff2af2d89 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/source/KuduLookupFunction.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/source/KuduLookupFunction.java
@@ -49,6 +49,10 @@ import java.util.stream.Collectors;
 
 import static org.apache.flink.util.Preconditions.checkNotNull;
 import static org.apache.flink.util.TimeUtils.parseDuration;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_OPERATION_TIMEOUT_IN_MS;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_SOCKET_READ_TIMEOUT_IN_MS;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DISABLED_STATISTICS;
 import static org.apache.kudu.client.KuduPredicate.ComparisonOp.EQUAL;
 
 /**
@@ -129,12 +133,25 @@ public class KuduLookupFunction extends 
TableFunction<Row> {
                             .build();
         }
 
-        this.client = new KuduClient.KuduClientBuilder(masters).build();
+        this.client = buildKuduClient();
 
         this.table = client.openTable(tableName);
         LOG.info("KuduLookupFunction opened.");
     }
 
+    private KuduClient buildKuduClient() {
+        KuduClient.KuduClientBuilder builder = new 
KuduClient.KuduClientBuilder(masters);
+        if (configuration.getBoolean(DISABLED_STATISTICS)) {
+            builder.disableStatistics();
+        }
+        
builder.defaultAdminOperationTimeoutMs(configuration.getLong(DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS));
+        
builder.defaultOperationTimeoutMs(configuration.getLong(DEFAULT_OPERATION_TIMEOUT_IN_MS));
+        
builder.defaultSocketReadTimeoutMs(configuration.getLong(DEFAULT_SOCKET_READ_TIMEOUT_IN_MS));
+
+        return builder
+                .build();
+    }
+
     @Override
     public void close() throws Exception {
         super.close();
diff --git 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/table/KuduDynamicTableFactory.java
 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/table/KuduDynamicTableFactory.java
index 642b9b617..e7d855a45 100644
--- 
a/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/table/KuduDynamicTableFactory.java
+++ 
b/inlong-sort/sort-connectors/kudu/src/main/java/org/apache/inlong/sort/kudu/table/KuduDynamicTableFactory.java
@@ -44,15 +44,19 @@ import static 
org.apache.inlong.sort.base.Constants.INLONG_AUDIT;
 import static org.apache.inlong.sort.base.Constants.INLONG_METRIC;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.CONNECTOR_MASTERS;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.CONNECTOR_TABLE;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_OPERATION_TIMEOUT_IN_MS;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_SOCKET_READ_TIMEOUT_IN_MS;
 import static 
org.apache.inlong.sort.kudu.common.KuduOptions.ENABLE_KEY_FIELD_CHECK;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.FLUSH_MODE;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_BUFFER_SIZE;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_BUFFER_TIME;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_CACHE_SIZE;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_CACHE_TIME;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS;
 import static org.apache.inlong.sort.kudu.common.KuduOptions.MAX_RETRIES;
 import static 
org.apache.inlong.sort.kudu.common.KuduOptions.SINK_KEY_FIELD_NAMES;
 import static 
org.apache.inlong.sort.kudu.common.KuduOptions.SINK_START_NEW_CHAIN;
+import static 
org.apache.inlong.sort.kudu.common.KuduOptions.DISABLED_STATISTICS;
 import static 
org.apache.inlong.sort.kudu.common.KuduOptions.WRITE_THREAD_COUNT;
 import static 
org.apache.inlong.sort.kudu.common.KuduValidator.CONNECTOR_TYPE_VALUE_KUDU;
 
@@ -175,6 +179,10 @@ public class KuduDynamicTableFactory
         options.add(MAX_BUFFER_TIME);
         options.add(SINK_KEY_FIELD_NAMES);
         options.add(ENABLE_KEY_FIELD_CHECK);
+        options.add(DEFAULT_ADMIN_OPERATION_TIMEOUT_IN_MS);
+        options.add(DEFAULT_OPERATION_TIMEOUT_IN_MS);
+        options.add(DEFAULT_SOCKET_READ_TIMEOUT_IN_MS);
+        options.add(DISABLED_STATISTICS);
 
         options.add(INLONG_METRIC);
         options.add(INLONG_AUDIT);

Reply via email to