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);