This is an automated email from the ASF dual-hosted git repository.
fcsaky pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kudu.git
The following commit(s) were added to refs/heads/main by this push:
new eb05fd5 [FLINK-35114] Remove old Table API implementations, update
Table API and Schema stack
eb05fd5 is described below
commit eb05fd588ff529fa6babf25c23b0f9a03e4891d6
Author: Ferenc Csaky <[email protected]>
AuthorDate: Tue Dec 3 17:16:58 2024 +0100
[FLINK-35114] Remove old Table API implementations, update Table API and
Schema stack
---
.../connector/kudu/connector/KuduTableInfo.java | 2 +-
.../writer/AbstractSingleOperationMapper.java | 8 +
.../writer/RowDataUpsertOperationMapper.java | 11 +-
.../kudu/format/AbstractKuduInputFormat.java | 7 +-
.../connector/kudu/table/KuduCommonOptions.java | 33 +++
.../kudu/table/KuduDynamicTableFactory.java | 157 +++++++++++
.../kudu/table/KuduDynamicTableOptions.java | 100 +++++++
.../table/{dynamic => }/KuduDynamicTableSink.java | 16 +-
.../{dynamic => }/KuduDynamicTableSource.java | 140 ++++------
.../KuduCatalog.java} | 94 +++----
.../{dynamic => }/catalog/KuduCatalogFactory.java | 33 ++-
.../kudu/table/catalog/KuduCatalogOptions.java | 35 +++
.../dynamic/KuduDynamicTableSourceSinkFactory.java | 235 -----------------
.../table/function/lookup/KuduLookupOptions.java | 78 ------
.../function/lookup/KuduRowDataLookupFunction.java | 172 +++---------
.../connector/kudu/table/utils/KuduTableUtils.java | 77 ++----
.../org.apache.flink.table.factories.Factory | 3 +-
.../org.apache.flink.table.factories.TableFactory | 17 --
.../connector/kudu/connector/KuduTestBase.java | 18 +-
.../table/{dynamic => }/KuduDynamicSinkTest.java | 8 +-
.../table/{dynamic => }/KuduDynamicSourceTest.java | 4 +-
.../kudu/table/KuduDynamicTableFactoryTest.java | 233 ++++++++++++++++
.../KuduRowDataLookupFunctionTest.java | 60 +----
.../kudu/table/KuduTableSourceITCase.java | 87 ++++++
.../connector/kudu/table/KuduTableTestUtils.java | 44 ++++
.../kudu/table/catalog/KuduCatalogTest.java | 292 +++++++++++++++++++++
26 files changed, 1211 insertions(+), 753 deletions(-)
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/KuduTableInfo.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/KuduTableInfo.java
index f8ad4fa..e2706db 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/KuduTableInfo.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/KuduTableInfo.java
@@ -39,7 +39,7 @@ import static
org.apache.flink.util.Preconditions.checkNotNull;
@PublicEvolving
public class KuduTableInfo implements Serializable {
- private String name;
+ private final String name;
private CreateTableOptionsFactory createTableOptionsFactory = null;
private ColumnSchemasFactory schemasFactory = null;
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/AbstractSingleOperationMapper.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/AbstractSingleOperationMapper.java
index 9a29207..ae3ec50 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/AbstractSingleOperationMapper.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/AbstractSingleOperationMapper.java
@@ -44,10 +44,18 @@ public abstract class AbstractSingleOperationMapper<T>
implements KuduOperationM
protected final String[] columnNames;
private final KuduOperation operation;
+ protected AbstractSingleOperationMapper(List<String> columnNames) {
+ this(columnNames, null);
+ }
+
protected AbstractSingleOperationMapper(String[] columnNames) {
this(columnNames, null);
}
+ public AbstractSingleOperationMapper(List<String> columnNames,
KuduOperation operation) {
+ this(columnNames.toArray(new String[0]), operation);
+ }
+
public AbstractSingleOperationMapper(String[] columnNames, KuduOperation
operation) {
this.columnNames = columnNames;
this.operation = operation;
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/RowDataUpsertOperationMapper.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/RowDataUpsertOperationMapper.java
index 5b0f56a..5dfca36 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/RowDataUpsertOperationMapper.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/connector/writer/RowDataUpsertOperationMapper.java
@@ -18,7 +18,7 @@
package org.apache.flink.connector.kudu.connector.writer;
import org.apache.flink.annotation.Internal;
-import org.apache.flink.table.api.TableSchema;
+import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.data.DecimalData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.StringData;
@@ -31,7 +31,6 @@ import org.apache.kudu.client.Operation;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.Arrays;
import java.util.Optional;
import static
org.apache.flink.table.types.logical.utils.LogicalTypeChecks.getPrecision;
@@ -47,12 +46,12 @@ public class RowDataUpsertOperationMapper extends
AbstractSingleOperationMapper<
private static final int MIN_TIMESTAMP_PRECISION = 0;
private static final int MAX_TIMESTAMP_PRECISION = 6;
- private LogicalType[] logicalTypes;
+ private final LogicalType[] logicalTypes;
- public RowDataUpsertOperationMapper(TableSchema schema) {
- super(schema.getFieldNames());
+ public RowDataUpsertOperationMapper(ResolvedSchema schema) {
+ super(schema.getColumnNames());
logicalTypes =
- Arrays.stream(schema.getFieldDataTypes())
+ schema.getColumnDataTypes().stream()
.map(DataType::getLogicalType)
.toArray(LogicalType[]::new);
}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/format/AbstractKuduInputFormat.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/format/AbstractKuduInputFormat.java
index 157b5eb..641a583 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/format/AbstractKuduInputFormat.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/format/AbstractKuduInputFormat.java
@@ -30,7 +30,7 @@ import
org.apache.flink.connector.kudu.connector.reader.KuduInputSplit;
import org.apache.flink.connector.kudu.connector.reader.KuduReader;
import org.apache.flink.connector.kudu.connector.reader.KuduReaderConfig;
import org.apache.flink.connector.kudu.connector.reader.KuduReaderIterator;
-import
org.apache.flink.connector.kudu.table.dynamic.catalog.KuduDynamicCatalog;
+import org.apache.flink.connector.kudu.table.catalog.KuduCatalog;
import org.apache.flink.core.io.InputSplitAssigner;
import org.apache.kudu.client.KuduException;
@@ -48,9 +48,8 @@ import static
org.apache.flink.util.Preconditions.checkNotNull;
* KuduTableInfo}) in both batch and stream programs. Rows of the Kudu table
are mapped to {@link T}
* instances that can converted to other data types by the user later if
necessary.
*
- * <p>For programmatic access to the schema of the input rows users can use
the {@link
- * KuduDynamicCatalog} or overwrite the column order manually by providing a
list of projected
- * column names.
+ * <p>For programmatic access to the schema of the input rows users can use
the {@link KuduCatalog}
+ * or overwrite the column order manually by providing a list of projected
column names.
*/
@PublicEvolving
public abstract class AbstractKuduInputFormat<T> extends RichInputFormat<T,
KuduInputSplit>
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduCommonOptions.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduCommonOptions.java
new file mode 100644
index 0000000..52a66cd
--- /dev/null
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduCommonOptions.java
@@ -0,0 +1,33 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.annotation.PublicEvolving;
+import org.apache.flink.configuration.ConfigOption;
+import org.apache.flink.configuration.ConfigOptions;
+
+/** Common options for Kudu tables and catalogs. */
+@PublicEvolving
+public class KuduCommonOptions {
+
+ public static final ConfigOption<String> KUDU_MASTERS =
+ ConfigOptions.key("kudu.masters")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("kudu's master server address");
+}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactory.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactory.java
new file mode 100644
index 0000000..a03cc2c
--- /dev/null
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactory.java
@@ -0,0 +1,157 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.configuration.ConfigOption;
+import org.apache.flink.configuration.ReadableConfig;
+import org.apache.flink.connector.kudu.connector.KuduTableInfo;
+import org.apache.flink.connector.kudu.connector.reader.KuduReaderConfig;
+import org.apache.flink.connector.kudu.connector.writer.KuduWriterConfig;
+import org.apache.flink.connector.kudu.table.utils.KuduTableUtils;
+import org.apache.flink.table.catalog.ResolvedSchema;
+import org.apache.flink.table.connector.sink.DynamicTableSink;
+import org.apache.flink.table.connector.source.DynamicTableSource;
+import org.apache.flink.table.connector.source.lookup.LookupOptions;
+import org.apache.flink.table.connector.source.lookup.cache.DefaultLookupCache;
+import org.apache.flink.table.connector.source.lookup.cache.LookupCache;
+import org.apache.flink.table.factories.DynamicTableSinkFactory;
+import org.apache.flink.table.factories.DynamicTableSourceFactory;
+import org.apache.flink.table.factories.FactoryUtil;
+
+import org.apache.kudu.shaded.com.google.common.collect.Sets;
+
+import javax.annotation.Nullable;
+
+import java.util.Set;
+
+import static
org.apache.flink.connector.kudu.table.KuduCommonOptions.KUDU_MASTERS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.IDENTIFIER;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_FLUSH_INTERVAL;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_PARTITION_NUMS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_IGNORE_DUPLICATE;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_IGNORE_NOT_FOUND;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_MAX_BUFFER_SIZE;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_OPERATION_TIMEOUT;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_PRIMARY_KEY_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_REPLICAS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_SCAN_ROW_SIZE;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_TABLE;
+
+/**
+ * Factory for creating configured instances of {@link
KuduDynamicTableSource}/{@link
+ * KuduDynamicTableSink} in a stream environment.
+ */
+public class KuduDynamicTableFactory implements DynamicTableSourceFactory,
DynamicTableSinkFactory {
+
+ @Override
+ public String factoryIdentifier() {
+ return IDENTIFIER;
+ }
+
+ @Override
+ public Set<ConfigOption<?>> requiredOptions() {
+ return Sets.newHashSet(KUDU_MASTERS);
+ }
+
+ @Override
+ public Set<ConfigOption<?>> optionalOptions() {
+ return Sets.newHashSet(
+ KUDU_TABLE,
+ KUDU_HASH_COLS,
+ KUDU_HASH_PARTITION_NUMS,
+ KUDU_PRIMARY_KEY_COLS,
+ KUDU_SCAN_ROW_SIZE,
+ KUDU_REPLICAS,
+ KUDU_MAX_BUFFER_SIZE,
+ KUDU_MAX_BUFFER_SIZE,
+ KUDU_OPERATION_TIMEOUT,
+ KUDU_FLUSH_INTERVAL,
+ KUDU_IGNORE_NOT_FOUND,
+ KUDU_IGNORE_DUPLICATE,
+ LookupOptions.CACHE_TYPE,
+ LookupOptions.PARTIAL_CACHE_MAX_ROWS,
+ LookupOptions.PARTIAL_CACHE_EXPIRE_AFTER_ACCESS,
+ LookupOptions.PARTIAL_CACHE_EXPIRE_AFTER_WRITE,
+ LookupOptions.PARTIAL_CACHE_CACHE_MISSING_KEY,
+ LookupOptions.MAX_RETRIES);
+ }
+
+ @Override
+ public DynamicTableSink createDynamicTableSink(Context context) {
+ final ReadableConfig config = getValidatedConfig(context);
+
+ final String tableName =
+ config.getOptional(KUDU_TABLE)
+ .orElse(context.getObjectIdentifier().getObjectName());
+ final ResolvedSchema schema =
context.getCatalogTable().getResolvedSchema();
+ final KuduTableInfo tableInfo =
+ KuduTableUtils.createTableInfo(
+ tableName, schema,
context.getCatalogTable().toProperties());
+
+ final KuduWriterConfig.Builder configBuilder =
+ KuduWriterConfig.Builder.setMasters(config.get(KUDU_MASTERS))
+
.setOperationTimeout(config.get(KUDU_OPERATION_TIMEOUT).toMillis())
+ .setFlushInterval((int)
config.get(KUDU_FLUSH_INTERVAL).toMillis())
+ .setMaxBufferSize(config.get(KUDU_MAX_BUFFER_SIZE))
+ .setIgnoreNotFound(config.get(KUDU_IGNORE_NOT_FOUND))
+ .setIgnoreDuplicate(config.get(KUDU_IGNORE_DUPLICATE));
+
+ return new KuduDynamicTableSink(configBuilder, tableInfo, schema);
+ }
+
+ @Override
+ public DynamicTableSource createDynamicTableSource(Context context) {
+ final ReadableConfig config = getValidatedConfig(context);
+
+ final String tableName =
+ config.getOptional(KUDU_TABLE)
+ .orElse(context.getObjectIdentifier().getObjectName());
+ final KuduTableInfo tableInfo =
+ KuduTableUtils.createTableInfo(
+ tableName,
+ context.getCatalogTable().getResolvedSchema(),
+ context.getCatalogTable().toProperties());
+
+ final KuduReaderConfig.Builder readerConfigBuilder =
+ KuduReaderConfig.Builder.setMasters(config.get(KUDU_MASTERS))
+ .setRowLimit(config.get(KUDU_SCAN_ROW_SIZE));
+
+ return new KuduDynamicTableSource(
+ readerConfigBuilder,
+ tableInfo,
+ context.getPhysicalRowDataType(),
+ config.get(LookupOptions.MAX_RETRIES),
+ getLookupCache(config));
+ }
+
+ private ReadableConfig getValidatedConfig(Context context) {
+ final FactoryUtil.TableFactoryHelper helper =
+ FactoryUtil.createTableFactoryHelper(this, context);
+ helper.validate();
+
+ return helper.getOptions();
+ }
+
+ @Nullable
+ private LookupCache getLookupCache(ReadableConfig config) {
+ return
LookupOptions.LookupCacheType.PARTIAL.equals(config.get(LookupOptions.CACHE_TYPE))
+ ? DefaultLookupCache.fromConfig(config)
+ : null;
+ }
+}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableOptions.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableOptions.java
new file mode 100644
index 0000000..9b4f11b
--- /dev/null
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableOptions.java
@@ -0,0 +1,100 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.annotation.PublicEvolving;
+import org.apache.flink.configuration.ConfigOption;
+import org.apache.flink.configuration.ConfigOptions;
+
+import java.time.Duration;
+
+/** Kudu table options. */
+@PublicEvolving
+public class KuduDynamicTableOptions {
+
+ public static final String IDENTIFIER = "kudu";
+
+ public static final ConfigOption<String> KUDU_TABLE =
+ ConfigOptions.key("kudu.table")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("kudu's table name");
+
+ public static final ConfigOption<String> KUDU_HASH_COLS =
+ ConfigOptions.key("kudu.hash-columns")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("kudu's hash columns");
+
+ public static final ConfigOption<Integer> KUDU_REPLICAS =
+ ConfigOptions.key("kudu.replicas")
+ .intType()
+ .defaultValue(3)
+ .withDescription("kudu's replica nums");
+
+ public static final ConfigOption<Integer> KUDU_MAX_BUFFER_SIZE =
+ ConfigOptions.key("kudu.max-buffer-size")
+ .intType()
+ .defaultValue(1000)
+ .withDescription("kudu's max buffer size");
+
+ public static final ConfigOption<Duration> KUDU_FLUSH_INTERVAL =
+ ConfigOptions.key("kudu.flush-interval")
+ .durationType()
+ .defaultValue(Duration.ofMillis(1000))
+ .withDescription("kudu's data flush interval");
+
+ public static final ConfigOption<Duration> KUDU_OPERATION_TIMEOUT =
+ ConfigOptions.key("kudu.operation-timeout")
+ .durationType()
+ .defaultValue(Duration.ofSeconds(30))
+ .withDescription("kudu's operation timeout");
+
+ public static final ConfigOption<Boolean> KUDU_IGNORE_NOT_FOUND =
+ ConfigOptions.key("kudu.ignore-not-found")
+ .booleanType()
+ .defaultValue(false)
+ .withDescription("if true, ignore all not found rows");
+
+ public static final ConfigOption<Boolean> KUDU_IGNORE_DUPLICATE =
+ ConfigOptions.key("kudu.ignore-duplicate")
+ .booleanType()
+ .defaultValue(false)
+ .withDescription("if true, ignore all duplicate rows");
+
+ public static final ConfigOption<Integer> KUDU_HASH_PARTITION_NUMS =
+ ConfigOptions.key("kudu.hash-partition-nums")
+ .intType()
+ .defaultValue(KUDU_REPLICAS.defaultValue() * 2)
+ .withDescription(
+ "kudu's hash partition bucket nums, defaultValue
is 2 * replica nums");
+
+ public static final ConfigOption<String> KUDU_PRIMARY_KEY_COLS =
+ ConfigOptions.key("kudu.primary-key-columns")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("kudu's primary key, primary key must be
ordered");
+
+ public static final ConfigOption<Integer> KUDU_SCAN_ROW_SIZE =
+ ConfigOptions.key("kudu.scan.row-size")
+ .intType()
+ .defaultValue(0)
+ .withDescription("kudu's scan row size");
+
+ private KuduDynamicTableOptions() {}
+}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSink.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSink.java
similarity index 90%
rename from
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSink.java
rename to
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSink.java
index 491baa6..870c364 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSink.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSink.java
@@ -15,13 +15,13 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic;
+package org.apache.flink.connector.kudu.table;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import org.apache.flink.connector.kudu.connector.writer.KuduWriterConfig;
import
org.apache.flink.connector.kudu.connector.writer.RowDataUpsertOperationMapper;
import org.apache.flink.connector.kudu.sink.KuduSink;
-import org.apache.flink.table.api.TableSchema;
+import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.connector.ChangelogMode;
import org.apache.flink.table.connector.sink.DynamicTableSink;
import org.apache.flink.table.connector.sink.SinkV2Provider;
@@ -34,16 +34,16 @@ import java.util.Objects;
/** A {@link KuduDynamicTableSink} for Kudu. */
public class KuduDynamicTableSink implements DynamicTableSink {
private final KuduWriterConfig.Builder writerConfigBuilder;
- private final TableSchema flinkSchema;
private final KuduTableInfo tableInfo;
+ private final ResolvedSchema flinkSchema;
public KuduDynamicTableSink(
KuduWriterConfig.Builder writerConfigBuilder,
- TableSchema flinkSchema,
- KuduTableInfo tableInfo) {
+ KuduTableInfo tableInfo,
+ ResolvedSchema flinkSchema) {
this.writerConfigBuilder = writerConfigBuilder;
- this.flinkSchema = flinkSchema;
this.tableInfo = tableInfo;
+ this.flinkSchema = flinkSchema;
}
@Override
@@ -59,7 +59,7 @@ public class KuduDynamicTableSink implements DynamicTableSink
{
private void validatePrimaryKey(ChangelogMode requestedMode) {
Preconditions.checkState(
ChangelogMode.insertOnly().equals(requestedMode)
- ||
this.tableInfo.getSchema().getPrimaryKeyColumnCount() != 0,
+ || tableInfo.getSchema().getPrimaryKeyColumnCount() !=
0,
"please declare primary key for sink table when query contains
update/delete record.");
}
@@ -76,7 +76,7 @@ public class KuduDynamicTableSink implements DynamicTableSink
{
@Override
public DynamicTableSink copy() {
- return new KuduDynamicTableSink(this.writerConfigBuilder,
this.flinkSchema, this.tableInfo);
+ return new KuduDynamicTableSink(writerConfigBuilder, tableInfo,
flinkSchema);
}
@Override
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSource.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSource.java
similarity index 53%
rename from
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSource.java
rename to
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSource.java
index bc19696..92adc3a 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSource.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/KuduDynamicTableSource.java
@@ -15,45 +15,41 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic;
+package org.apache.flink.connector.kudu.table;
import org.apache.flink.connector.kudu.connector.KuduFilterInfo;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import
org.apache.flink.connector.kudu.connector.converter.RowResultRowDataConverter;
import org.apache.flink.connector.kudu.connector.reader.KuduReaderConfig;
import org.apache.flink.connector.kudu.format.KuduRowDataInputFormat;
-import org.apache.flink.connector.kudu.table.function.lookup.KuduLookupOptions;
import
org.apache.flink.connector.kudu.table.function.lookup.KuduRowDataLookupFunction;
import org.apache.flink.connector.kudu.table.utils.KuduTableUtils;
-import org.apache.flink.table.api.TableSchema;
import org.apache.flink.table.connector.ChangelogMode;
+import org.apache.flink.table.connector.Projection;
import org.apache.flink.table.connector.source.DynamicTableSource;
import org.apache.flink.table.connector.source.InputFormatProvider;
import org.apache.flink.table.connector.source.LookupTableSource;
import org.apache.flink.table.connector.source.ScanTableSource;
-import org.apache.flink.table.connector.source.TableFunctionProvider;
import
org.apache.flink.table.connector.source.abilities.SupportsFilterPushDown;
import org.apache.flink.table.connector.source.abilities.SupportsLimitPushDown;
import
org.apache.flink.table.connector.source.abilities.SupportsProjectionPushDown;
+import org.apache.flink.table.connector.source.lookup.LookupFunctionProvider;
+import
org.apache.flink.table.connector.source.lookup.PartialCachingLookupProvider;
+import org.apache.flink.table.connector.source.lookup.cache.LookupCache;
import org.apache.flink.table.expressions.ResolvedExpression;
import org.apache.flink.table.types.DataType;
-import org.apache.flink.table.types.FieldsDataType;
-import org.apache.flink.table.types.logical.RowType;
-import org.apache.flink.table.types.utils.DataTypeUtils;
-import org.apache.flink.util.Preconditions;
import org.apache.commons.collections.CollectionUtils;
-import org.apache.kudu.shaded.com.google.common.collect.Lists;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import java.util.Arrays;
+import javax.annotation.Nullable;
+
+import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
-import static
org.apache.flink.table.utils.TableSchemaUtils.containsPhysicalColumnsOnly;
+import static org.apache.flink.util.Preconditions.checkArgument;
/** A {@link DynamicTableSource} for Kudu. */
public class KuduDynamicTableSource
@@ -63,55 +59,48 @@ public class KuduDynamicTableSource
LookupTableSource,
SupportsFilterPushDown {
- private static final Logger LOG =
LoggerFactory.getLogger(KuduDynamicTableSource.class);
private final KuduTableInfo tableInfo;
- private final KuduLookupOptions kuduLookupOptions;
- private final KuduRowDataInputFormat kuduRowDataInputFormat;
- private final transient List<KuduFilterInfo> predicates =
Lists.newArrayList();
+ private final int lookupMaxRetryTimes;
+ @Nullable private final LookupCache cache;
+ private final transient List<KuduFilterInfo> predicates = new
ArrayList<>();
+
private KuduReaderConfig.Builder configBuilder;
- private TableSchema physicalSchema;
- private String[] projectedFields;
+ private DataType physicalRowDataType;
private transient List<ResolvedExpression> filters;
public KuduDynamicTableSource(
KuduReaderConfig.Builder configBuilder,
KuduTableInfo tableInfo,
- TableSchema physicalSchema,
- String[] projectedFields,
- KuduLookupOptions kuduLookupOptions) {
+ DataType physicalRowDataType,
+ int lookupMaxRetryTimes,
+ @Nullable LookupCache cache) {
this.configBuilder = configBuilder;
this.tableInfo = tableInfo;
- this.physicalSchema = physicalSchema;
- this.projectedFields = projectedFields;
- this.kuduRowDataInputFormat =
- new KuduRowDataInputFormat(
- configBuilder.build(),
- new RowResultRowDataConverter(),
- tableInfo,
- predicates,
- projectedFields == null ? null :
Lists.newArrayList(projectedFields));
- this.kuduLookupOptions = kuduLookupOptions;
+ this.physicalRowDataType = physicalRowDataType;
+ this.lookupMaxRetryTimes = lookupMaxRetryTimes;
+ this.cache = cache;
}
@Override
public LookupRuntimeProvider getLookupRuntimeProvider(LookupContext
context) {
- int keysLen = context.getKeys().length;
- String[] keyNames = new String[keysLen];
- for (int i = 0; i < keyNames.length; ++i) {
+ String[] keyNames = new String[context.getKeys().length];
+ for (int i = 0; i < keyNames.length; i++) {
int[] innerKeyArr = context.getKeys()[i];
- Preconditions.checkArgument(
- innerKeyArr.length == 1, "Kudu only support non-nested
look up keys");
- keyNames[i] = this.physicalSchema.getFieldNames()[innerKeyArr[0]];
+ checkArgument(innerKeyArr.length == 1, "Kudu only supports
non-nested lookup keys");
+ keyNames[i] =
DataType.getFieldNames(physicalRowDataType).get(innerKeyArr[0]);
}
- KuduRowDataLookupFunction rowDataLookupFunction =
- KuduRowDataLookupFunction.Builder.options()
- .keyNames(keyNames)
- .kuduReaderConfig(configBuilder.build())
- .projectedFields(projectedFields)
- .tableInfo(tableInfo)
- .kuduLookupOptions(kuduLookupOptions)
- .build();
- return TableFunctionProvider.of(rowDataLookupFunction);
+
+ KuduRowDataLookupFunction lookupFunction =
+ new KuduRowDataLookupFunction(
+ keyNames,
+ tableInfo,
+ configBuilder.build(),
+ DataType.getFieldNames(physicalRowDataType),
+ lookupMaxRetryTimes);
+
+ return cache == null
+ ? LookupFunctionProvider.of(lookupFunction)
+ : PartialCachingLookupProvider.of(lookupFunction, cache);
}
@Override
@@ -121,32 +110,30 @@ public class KuduDynamicTableSource
@Override
public ScanRuntimeProvider getScanRuntimeProvider(ScanContext
runtimeProviderContext) {
- if (CollectionUtils.isNotEmpty(this.filters)) {
- for (ResolvedExpression filter : this.filters) {
+ if (CollectionUtils.isNotEmpty(filters)) {
+ for (ResolvedExpression filter : filters) {
Optional<KuduFilterInfo> kuduFilterInfo =
KuduTableUtils.toKuduFilterInfo(filter);
if (kuduFilterInfo != null && kuduFilterInfo.isPresent()) {
- this.predicates.add(kuduFilterInfo.get());
+ predicates.add(kuduFilterInfo.get());
}
}
}
+
KuduRowDataInputFormat inputFormat =
new KuduRowDataInputFormat(
configBuilder.build(),
new RowResultRowDataConverter(),
tableInfo,
- this.predicates,
- projectedFields == null ? null :
Lists.newArrayList(projectedFields));
+ predicates,
+ DataType.getFieldNames(physicalRowDataType));
+
return InputFormatProvider.of(inputFormat);
}
@Override
public DynamicTableSource copy() {
return new KuduDynamicTableSource(
- this.configBuilder,
- this.tableInfo,
- this.physicalSchema,
- this.projectedFields,
- this.kuduLookupOptions);
+ configBuilder, tableInfo, physicalRowDataType,
lookupMaxRetryTimes, cache);
}
@Override
@@ -162,25 +149,7 @@ public class KuduDynamicTableSource
@Override
public void applyProjection(int[][] projectedFields, DataType
producedDataType) {
- // parser projectFields
- this.physicalSchema = projectSchema(this.physicalSchema,
projectedFields);
- this.projectedFields = physicalSchema.getFieldNames();
- }
-
- private TableSchema projectSchema(TableSchema tableSchema, int[][]
projectedFields) {
- Preconditions.checkArgument(
- containsPhysicalColumnsOnly(tableSchema),
- "Projection is only supported for physical columns.");
- TableSchema.Builder builder = TableSchema.builder();
-
- FieldsDataType fields =
- (FieldsDataType)
- DataTypeUtils.projectRow(tableSchema.toRowDataType(),
projectedFields);
- RowType topFields = (RowType) fields.getLogicalType();
- for (int i = 0; i < topFields.getFieldCount(); i++) {
- builder.field(topFields.getFieldNames().get(i),
fields.getChildren().get(i));
- }
- return builder.build();
+ physicalRowDataType =
Projection.of(projectedFields).project(physicalRowDataType);
}
@Override
@@ -194,27 +163,16 @@ public class KuduDynamicTableSource
KuduDynamicTableSource that = (KuduDynamicTableSource) o;
return Objects.equals(configBuilder, that.configBuilder)
&& Objects.equals(tableInfo, that.tableInfo)
- && Objects.equals(physicalSchema, that.physicalSchema)
- && Arrays.equals(projectedFields, that.projectedFields)
- && Objects.equals(kuduLookupOptions, that.kuduLookupOptions)
- && Objects.equals(kuduRowDataInputFormat,
that.kuduRowDataInputFormat)
+ && Objects.equals(cache, that.cache)
+ && Objects.equals(physicalRowDataType,
that.physicalRowDataType)
&& Objects.equals(filters, that.filters)
&& Objects.equals(predicates, that.predicates);
}
@Override
public int hashCode() {
- int result =
- Objects.hash(
- configBuilder,
- tableInfo,
- physicalSchema,
- kuduLookupOptions,
- kuduRowDataInputFormat,
- filters,
- predicates);
- result = 31 * result + Arrays.hashCode(projectedFields);
- return result;
+ return Objects.hash(
+ configBuilder, tableInfo, cache, physicalRowDataType, filters,
predicates);
}
@Override
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduDynamicCatalog.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalog.java
similarity index 81%
rename from
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduDynamicCatalog.java
rename to
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalog.java
index 5f51e06..34b7cef 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduDynamicCatalog.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalog.java
@@ -16,21 +16,21 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic.catalog;
+package org.apache.flink.connector.kudu.table.catalog;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import org.apache.flink.connector.kudu.table.AbstractReadOnlyCatalog;
-import
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory;
+import org.apache.flink.connector.kudu.table.KuduDynamicTableFactory;
import org.apache.flink.connector.kudu.table.utils.KuduTableUtils;
-import org.apache.flink.table.api.TableSchema;
import org.apache.flink.table.catalog.CatalogBaseTable;
import org.apache.flink.table.catalog.CatalogDatabase;
import org.apache.flink.table.catalog.CatalogDatabaseImpl;
import org.apache.flink.table.catalog.CatalogFunction;
import org.apache.flink.table.catalog.CatalogPartitionSpec;
import org.apache.flink.table.catalog.CatalogTable;
-import org.apache.flink.table.catalog.CatalogTableImpl;
import org.apache.flink.table.catalog.ObjectPath;
+import org.apache.flink.table.catalog.ResolvedCatalogBaseTable;
+import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.catalog.exceptions.CatalogException;
import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException;
import org.apache.flink.table.catalog.exceptions.FunctionNotExistException;
@@ -47,12 +47,12 @@ import org.apache.kudu.client.AlterTableOptions;
import org.apache.kudu.client.KuduClient;
import org.apache.kudu.client.KuduException;
import org.apache.kudu.client.KuduTable;
+import org.apache.kudu.shaded.com.google.common.collect.ImmutableSet;
import org.apache.kudu.shaded.com.google.common.collect.Lists;
import org.apache.kudu.shaded.com.google.common.collect.Sets;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
@@ -62,42 +62,44 @@ import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.KUDU_HASH_COLS;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.KUDU_HASH_PARTITION_NUMS;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.KUDU_PRIMARY_KEY_COLS;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.KUDU_REPLICAS;
+import static
org.apache.flink.connector.kudu.table.KuduCommonOptions.KUDU_MASTERS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_PARTITION_NUMS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_PRIMARY_KEY_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_REPLICAS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_TABLE;
import static org.apache.flink.util.Preconditions.checkArgument;
import static org.apache.flink.util.Preconditions.checkNotNull;
/** Catalog for reading and creating Kudu tables. */
-public class KuduDynamicCatalog extends AbstractReadOnlyCatalog {
+public class KuduCatalog extends AbstractReadOnlyCatalog {
- private static final Logger LOG =
LoggerFactory.getLogger(KuduDynamicCatalog.class);
- private final KuduDynamicTableSourceSinkFactory tableFactory =
- new KuduDynamicTableSourceSinkFactory();
+ private static final Logger LOG =
LoggerFactory.getLogger(KuduCatalog.class);
+ private final KuduDynamicTableFactory tableFactory = new
KuduDynamicTableFactory();
private final String kuduMasters;
private final KuduClient kuduClient;
/**
- * Create a new {@link KuduDynamicCatalog} with the specified catalog name
and kudu master
- * addresses.
+ * Create a new {@link KuduCatalog} with the specified kudu master
addresses.
*
- * @param catalogName Name of the catalog (used by the table environment)
* @param kuduMasters Connection address to Kudu
*/
- public KuduDynamicCatalog(String catalogName, String kuduMasters) {
- super(catalogName, "default_database");
- this.kuduMasters = kuduMasters;
- this.kuduClient = createClient();
+ public KuduCatalog(String kuduMasters) {
+ this("kudu", "default", kuduMasters);
}
/**
- * Create a new {@link KuduDynamicCatalog} with the specified kudu master
addresses.
+ * Create a new {@link KuduCatalog} with the specified catalog name,
default database, and kudu
+ * master addresses.
*
+ * @param catalogName Name of the catalog (used by the table environment)
+ * @param defaultDatabase Default Kudu database name
* @param kuduMasters Connection address to Kudu
*/
- public KuduDynamicCatalog(String kuduMasters) {
- this("kudu", kuduMasters);
+ public KuduCatalog(String catalogName, String defaultDatabase, String
kuduMasters) {
+ super(catalogName, defaultDatabase);
+ this.kuduMasters = kuduMasters;
+ this.kuduClient = createClient();
}
@Override
@@ -105,7 +107,7 @@ public class KuduDynamicCatalog extends
AbstractReadOnlyCatalog {
return Optional.of(getKuduTableFactory());
}
- public KuduDynamicTableSourceSinkFactory getKuduTableFactory() {
+ public KuduDynamicTableFactory getKuduTableFactory() {
return tableFactory;
}
@@ -171,15 +173,11 @@ public class KuduDynamicCatalog extends
AbstractReadOnlyCatalog {
try {
KuduTable kuduTable = kuduClient.openTable(tableName);
- // fixme base on TableSchema, TableSchema needs to be upgraded to
ResolvedSchema
- CatalogTableImpl table =
- new CatalogTableImpl(
-
KuduTableUtils.kuduToFlinkSchema(kuduTable.getSchema()),
- createTableProperties(
- tableName,
kuduTable.getSchema().getPrimaryKeyColumns()),
- tableName);
-
- return table;
+ return CatalogTable.of(
+ KuduTableUtils.kuduToFlinkSchema(kuduTable.getSchema()),
+ null,
+ Collections.emptyList(),
+ createTableProperties(tableName,
kuduTable.getSchema().getPrimaryKeyColumns()));
} catch (KuduException e) {
throw new CatalogException(e);
}
@@ -188,13 +186,13 @@ public class KuduDynamicCatalog extends
AbstractReadOnlyCatalog {
protected Map<String, String> createTableProperties(
String tableName, List<ColumnSchema> primaryKeyColumns) {
Map<String, String> props = new HashMap<>();
- props.put(KuduDynamicTableSourceSinkFactory.KUDU_MASTERS.key(),
kuduMasters);
+ props.put(KUDU_MASTERS.key(), kuduMasters);
+ props.put(KUDU_TABLE.key(), tableName);
String primaryKeyNames =
primaryKeyColumns.stream()
.map(ColumnSchema::getName)
.collect(Collectors.joining(","));
props.put(KUDU_PRIMARY_KEY_COLS.key(), primaryKeyNames);
- props.put(KuduDynamicTableSourceSinkFactory.KUDU_TABLE.key(),
tableName);
return props;
}
@@ -230,6 +228,8 @@ public class KuduDynamicCatalog extends
AbstractReadOnlyCatalog {
public void createTable(KuduTableInfo tableInfo, boolean ignoreIfExists)
throws CatalogException, TableAlreadyExistException {
+ checkNotNull(tableInfo);
+
ObjectPath path = getObjectPath(tableInfo.getName());
if (tableExists(path)) {
if (ignoreIfExists) {
@@ -250,38 +250,38 @@ public class KuduDynamicCatalog extends
AbstractReadOnlyCatalog {
@Override
public void createTable(ObjectPath tablePath, CatalogBaseTable table,
boolean ignoreIfExists)
throws TableAlreadyExistException {
+ checkNotNull(tablePath, "Table path must be provided.");
+ checkNotNull(table, "Table must be provided.");
+ checkArgument(table instanceof ResolvedCatalogBaseTable, "Table must
be resolved.");
+
Map<String, String> tableProperties = table.getOptions();
- TableSchema tableSchema = table.getSchema();
+ ResolvedSchema schema = ((ResolvedCatalogBaseTable<?>)
table).getResolvedSchema();
Set<String> optionalProperties =
- new HashSet<>(
- Arrays.asList(
- KUDU_REPLICAS.key(),
- KUDU_HASH_PARTITION_NUMS.key(),
- KUDU_HASH_COLS.key()));
- Set<String> requiredProperties = new HashSet<>();
+ ImmutableSet.of(
+ KUDU_REPLICAS.key(), KUDU_HASH_PARTITION_NUMS.key(),
KUDU_HASH_COLS.key());
- if (!tableSchema.getPrimaryKey().isPresent()) {
+ Set<String> requiredProperties = new HashSet<>();
+ if (!schema.getPrimaryKey().isPresent()) {
requiredProperties.add(KUDU_PRIMARY_KEY_COLS.key());
}
if (!tableProperties.keySet().containsAll(requiredProperties)) {
throw new CatalogException(
"Missing required property. The following properties must
be provided: "
- + requiredProperties.toString());
+ + requiredProperties);
}
Set<String> permittedProperties = Sets.union(requiredProperties,
optionalProperties);
if (!permittedProperties.containsAll(tableProperties.keySet())) {
throw new CatalogException(
"Unpermitted properties were given. The following
properties are allowed:"
- + permittedProperties.toString());
+ + permittedProperties);
}
String tableName = tablePath.getObjectName();
-
KuduTableInfo tableInfo =
- KuduTableUtils.createTableInfo(tableName, tableSchema,
tableProperties);
+ KuduTableUtils.createTableInfo(tableName, schema,
tableProperties);
createTable(tableInfo, ignoreIfExists);
}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduCatalogFactory.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogFactory.java
similarity index 66%
rename from
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduCatalogFactory.java
rename to
flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogFactory.java
index 21f5e4c..455a6b0 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/catalog/KuduCatalogFactory.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogFactory.java
@@ -16,7 +16,7 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic.catalog;
+package org.apache.flink.connector.kudu.table.catalog;
import org.apache.flink.annotation.Internal;
import org.apache.flink.configuration.ConfigOption;
@@ -24,39 +24,32 @@ import org.apache.flink.table.catalog.Catalog;
import org.apache.flink.table.factories.CatalogFactory;
import org.apache.flink.table.factories.FactoryUtil;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import org.apache.kudu.shaded.com.google.common.collect.Sets;
-import java.util.HashSet;
import java.util.Set;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.IDENTIFIER;
-import static
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory.KUDU_MASTERS;
+import static
org.apache.flink.connector.kudu.table.KuduCommonOptions.KUDU_MASTERS;
+import static
org.apache.flink.connector.kudu.table.catalog.KuduCatalogOptions.DEFAULT_DATABASE;
+import static
org.apache.flink.connector.kudu.table.catalog.KuduCatalogOptions.IDENTIFIER;
import static org.apache.flink.table.factories.FactoryUtil.PROPERTY_VERSION;
-/** Factory for {@link KuduDynamicCatalog}. */
+/** Factory for {@link KuduCatalog}. */
@Internal
public class KuduCatalogFactory implements CatalogFactory {
- private static final Logger LOG =
LoggerFactory.getLogger(KuduCatalogFactory.class);
-
@Override
public String factoryIdentifier() {
return IDENTIFIER;
}
@Override
- public Set<ConfigOption<?>> optionalOptions() {
- final Set<ConfigOption<?>> options = new HashSet<>();
- options.add(PROPERTY_VERSION);
- return options;
+ public Set<ConfigOption<?>> requiredOptions() {
+ return Sets.newHashSet(KUDU_MASTERS);
}
@Override
- public Set<ConfigOption<?>> requiredOptions() {
- final Set<ConfigOption<?>> options = new HashSet<>();
- options.add(KUDU_MASTERS);
- return options;
+ public Set<ConfigOption<?>> optionalOptions() {
+ return Sets.newHashSet(DEFAULT_DATABASE, PROPERTY_VERSION);
}
@Override
@@ -64,6 +57,10 @@ public class KuduCatalogFactory implements CatalogFactory {
final FactoryUtil.CatalogFactoryHelper helper =
FactoryUtil.createCatalogFactoryHelper(this, context);
helper.validate();
- return new KuduDynamicCatalog(context.getName(),
helper.getOptions().get(KUDU_MASTERS));
+
+ return new KuduCatalog(
+ context.getName(),
+ helper.getOptions().get(DEFAULT_DATABASE),
+ helper.getOptions().get(KUDU_MASTERS));
}
}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogOptions.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogOptions.java
new file mode 100644
index 0000000..9224913
--- /dev/null
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogOptions.java
@@ -0,0 +1,35 @@
+/*
+ * 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.flink.connector.kudu.table.catalog;
+
+import org.apache.flink.annotation.PublicEvolving;
+import org.apache.flink.configuration.ConfigOption;
+import org.apache.flink.configuration.ConfigOptions;
+import org.apache.flink.table.catalog.CommonCatalogOptions;
+
+/** Kudu catalog options. */
+@PublicEvolving
+public class KuduCatalogOptions {
+
+ public static final String IDENTIFIER = "kudu";
+
+ public static final ConfigOption<String> DEFAULT_DATABASE =
+ ConfigOptions.key(CommonCatalogOptions.DEFAULT_DATABASE_KEY)
+ .stringType()
+ .defaultValue("default");
+}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSourceSinkFactory.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSourceSinkFactory.java
deleted file mode 100644
index 9bf3816..0000000
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicTableSourceSinkFactory.java
+++ /dev/null
@@ -1,235 +0,0 @@
-/*
- * 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.flink.connector.kudu.table.dynamic;
-
-import org.apache.flink.configuration.ConfigOption;
-import org.apache.flink.configuration.ConfigOptions;
-import org.apache.flink.configuration.ReadableConfig;
-import org.apache.flink.connector.kudu.connector.KuduTableInfo;
-import org.apache.flink.connector.kudu.connector.reader.KuduReaderConfig;
-import org.apache.flink.connector.kudu.connector.writer.KuduWriterConfig;
-import org.apache.flink.connector.kudu.table.function.lookup.KuduLookupOptions;
-import org.apache.flink.connector.kudu.table.utils.KuduTableUtils;
-import org.apache.flink.table.api.TableSchema;
-import org.apache.flink.table.connector.sink.DynamicTableSink;
-import org.apache.flink.table.connector.source.DynamicTableSource;
-import org.apache.flink.table.factories.DynamicTableSinkFactory;
-import org.apache.flink.table.factories.DynamicTableSourceFactory;
-import org.apache.flink.table.factories.FactoryUtil;
-
-import org.apache.kudu.shaded.com.google.common.collect.Sets;
-
-import java.util.Optional;
-import java.util.Set;
-
-/**
- * Factory for creating configured instances of {@link
KuduDynamicTableSource}/{@link
- * KuduDynamicTableSink} in a stream environment.
- */
-public class KuduDynamicTableSourceSinkFactory
- implements DynamicTableSourceFactory, DynamicTableSinkFactory {
- public static final String IDENTIFIER = "kudu";
- public static final ConfigOption<String> KUDU_TABLE =
- ConfigOptions.key("kudu.table")
- .stringType()
- .noDefaultValue()
- .withDescription("kudu's table name");
-
- public static final ConfigOption<String> KUDU_MASTERS =
- ConfigOptions.key("kudu.masters")
- .stringType()
- .noDefaultValue()
- .withDescription("kudu's master server address");
-
- public static final ConfigOption<String> KUDU_HASH_COLS =
- ConfigOptions.key("kudu.hash-columns")
- .stringType()
- .noDefaultValue()
- .withDescription("kudu's hash columns");
-
- public static final ConfigOption<Integer> KUDU_REPLICAS =
- ConfigOptions.key("kudu.replicas")
- .intType()
- .defaultValue(3)
- .withDescription("kudu's replica nums");
-
- public static final ConfigOption<Integer> KUDU_MAX_BUFFER_SIZE =
- ConfigOptions.key("kudu.max-buffer-size")
- .intType()
- .noDefaultValue()
- .withDescription("kudu's max buffer size");
-
- public static final ConfigOption<Integer> KUDU_FLUSH_INTERVAL =
- ConfigOptions.key("kudu.flush-interval")
- .intType()
- .noDefaultValue()
- .withDescription("kudu's data flush interval");
-
- public static final ConfigOption<Long> KUDU_OPERATION_TIMEOUT =
- ConfigOptions.key("kudu.operation-timeout")
- .longType()
- .noDefaultValue()
- .withDescription("kudu's operation timeout");
-
- public static final ConfigOption<Boolean> KUDU_IGNORE_NOT_FOUND =
- ConfigOptions.key("kudu.ignore-not-found")
- .booleanType()
- .noDefaultValue()
- .withDescription("if true, ignore all not found rows");
-
- public static final ConfigOption<Boolean> KUDU_IGNORE_DUPLICATE =
- ConfigOptions.key("kudu.ignore-not-found")
- .booleanType()
- .noDefaultValue()
- .withDescription("if true, ignore all dulicate rows");
-
- public static final ConfigOption<Integer> KUDU_HASH_PARTITION_NUMS =
- ConfigOptions.key("kudu.hash-partition-nums")
- .intType()
- .defaultValue(KUDU_REPLICAS.defaultValue() * 2)
- .withDescription(
- "kudu's hash partition bucket nums, defaultValue
is 2 * replica nums");
-
- public static final ConfigOption<String> KUDU_PRIMARY_KEY_COLS =
- ConfigOptions.key("kudu.primary-key-columns")
- .stringType()
- .noDefaultValue()
- .withDescription("kudu's primary key, primary key must be
ordered");
-
- public static final ConfigOption<Integer> KUDU_SCAN_ROW_SIZE =
- ConfigOptions.key("kudu.scan.row-size")
- .intType()
- .defaultValue(0)
- .withDescription("kudu's scan row size");
-
- public static final ConfigOption<Long> KUDU_LOOKUP_CACHE_MAX_ROWS =
- ConfigOptions.key("kudu.lookup.cache.max-rows")
- .longType()
- .defaultValue(-1L)
- .withDescription(
- "the max number of rows of lookup cache, over this
value, the oldest rows will "
- + "be eliminated. \"cache.max-rows\" and
\"cache.ttl\" options must all be specified if any"
- + " of them is "
- + "specified. Cache is not enabled as
default.");
-
- public static final ConfigOption<Long> KUDU_LOOKUP_CACHE_TTL =
- ConfigOptions.key("kudu.lookup.cache.ttl")
- .longType()
- .defaultValue(-1L)
- .withDescription("the cache time to live.");
-
- public static final ConfigOption<Integer> KUDU_LOOKUP_MAX_RETRIES =
- ConfigOptions.key("kudu.lookup.max-retries")
- .intType()
- .defaultValue(3)
- .withDescription("the max retry times if lookup database
failed.");
-
- @Override
- public DynamicTableSink createDynamicTableSink(Context context) {
- ReadableConfig config = getReadableConfig(context);
- String masterAddresses = config.get(KUDU_MASTERS);
- String tableName = config.get(KUDU_TABLE);
- Optional<Long> operationTimeout =
config.getOptional(KUDU_OPERATION_TIMEOUT);
- Optional<Integer> flushInterval =
config.getOptional(KUDU_FLUSH_INTERVAL);
- Optional<Integer> bufferSize =
config.getOptional(KUDU_MAX_BUFFER_SIZE);
- Optional<Boolean> ignoreNotFound =
config.getOptional(KUDU_IGNORE_NOT_FOUND);
- Optional<Boolean> ignoreDuplicate =
config.getOptional(KUDU_IGNORE_DUPLICATE);
- TableSchema schema = context.getCatalogTable().getSchema();
- TableSchema physicalSchema =
KuduTableUtils.getSchemaWithSqlTimestamp(schema);
- KuduTableInfo tableInfo =
- KuduTableUtils.createTableInfo(
- tableName, schema,
context.getCatalogTable().toProperties());
-
- KuduWriterConfig.Builder configBuilder =
- KuduWriterConfig.Builder.setMasters(masterAddresses);
- operationTimeout.ifPresent(configBuilder::setOperationTimeout);
- flushInterval.ifPresent(configBuilder::setFlushInterval);
- bufferSize.ifPresent(configBuilder::setMaxBufferSize);
- ignoreNotFound.ifPresent(configBuilder::setIgnoreNotFound);
- ignoreDuplicate.ifPresent(configBuilder::setIgnoreDuplicate);
- return new KuduDynamicTableSink(configBuilder, physicalSchema,
tableInfo);
- }
-
- private ReadableConfig getReadableConfig(Context context) {
- FactoryUtil.TableFactoryHelper helper =
FactoryUtil.createTableFactoryHelper(this, context);
- return helper.getOptions();
- }
-
- @Override
- public DynamicTableSource createDynamicTableSource(Context context) {
- ReadableConfig config = getReadableConfig(context);
- String masterAddresses = config.get(KUDU_MASTERS);
-
- int scanRowSize = config.get(KUDU_SCAN_ROW_SIZE);
- long kuduCacheMaxRows = config.get(KUDU_LOOKUP_CACHE_MAX_ROWS);
- long kuduCacheTtl = config.get(KUDU_LOOKUP_CACHE_TTL);
- int kuduMaxReties = config.get(KUDU_LOOKUP_MAX_RETRIES);
-
- // build kudu lookup options
- KuduLookupOptions kuduLookupOptions =
- KuduLookupOptions.Builder.options()
- .withCacheMaxSize(kuduCacheMaxRows)
- .withCacheExpireMs(kuduCacheTtl)
- .withMaxRetryTimes(kuduMaxReties)
- .build();
-
- TableSchema schema = context.getCatalogTable().getSchema();
- TableSchema physicalSchema =
KuduTableUtils.getSchemaWithSqlTimestamp(schema);
- KuduTableInfo tableInfo =
- KuduTableUtils.createTableInfo(
- config.get(KUDU_TABLE), schema,
context.getCatalogTable().toProperties());
-
- KuduReaderConfig.Builder configBuilder =
-
KuduReaderConfig.Builder.setMasters(masterAddresses).setRowLimit(scanRowSize);
- return new KuduDynamicTableSource(
- configBuilder,
- tableInfo,
- physicalSchema,
- physicalSchema.getFieldNames(),
- kuduLookupOptions);
- }
-
- @Override
- public String factoryIdentifier() {
- return IDENTIFIER;
- }
-
- @Override
- public Set<ConfigOption<?>> requiredOptions() {
- return Sets.newHashSet(KUDU_TABLE, KUDU_MASTERS);
- }
-
- @Override
- public Set<ConfigOption<?>> optionalOptions() {
- return Sets.newHashSet(
- KUDU_HASH_COLS,
- KUDU_HASH_PARTITION_NUMS,
- KUDU_PRIMARY_KEY_COLS,
- KUDU_SCAN_ROW_SIZE,
- KUDU_REPLICAS,
- KUDU_MAX_BUFFER_SIZE,
- KUDU_MAX_BUFFER_SIZE,
- KUDU_OPERATION_TIMEOUT,
- KUDU_IGNORE_NOT_FOUND,
- KUDU_IGNORE_DUPLICATE,
- // lookup
- KUDU_LOOKUP_CACHE_MAX_ROWS,
- KUDU_LOOKUP_CACHE_TTL,
- KUDU_LOOKUP_MAX_RETRIES);
- }
-}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduLookupOptions.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduLookupOptions.java
deleted file mode 100644
index f05b988..0000000
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduLookupOptions.java
+++ /dev/null
@@ -1,78 +0,0 @@
-/*
- * 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.flink.connector.kudu.table.function.lookup;
-
-/** Options for the Kudu lookup. */
-public class KuduLookupOptions {
- private final long cacheMaxSize;
- private final long cacheExpireMs;
- private final int maxRetryTimes;
-
- public static Builder builder() {
- return new Builder();
- }
-
- public KuduLookupOptions(long cacheMaxSize, long cacheExpireMs, int
maxRetryTimes) {
- this.cacheMaxSize = cacheMaxSize;
- this.cacheExpireMs = cacheExpireMs;
- this.maxRetryTimes = maxRetryTimes;
- }
-
- public long getCacheMaxSize() {
- return cacheMaxSize;
- }
-
- public long getCacheExpireMs() {
- return cacheExpireMs;
- }
-
- public int getMaxRetryTimes() {
- return maxRetryTimes;
- }
-
- /** Builder for KuduLookupOptions. */
- public static final class Builder {
- private long cacheMaxSize;
- private long cacheExpireMs;
- private int maxRetryTimes;
-
- public static Builder options() {
- return new Builder();
- }
-
- public Builder withCacheMaxSize(long cacheMaxSize) {
- this.cacheMaxSize = cacheMaxSize;
- return this;
- }
-
- public Builder withCacheExpireMs(long cacheExpireMs) {
- this.cacheExpireMs = cacheExpireMs;
- return this;
- }
-
- public Builder withMaxRetryTimes(int maxRetryTimes) {
- this.maxRetryTimes = maxRetryTimes;
- return this;
- }
-
- public KuduLookupOptions build() {
- return new KuduLookupOptions(cacheMaxSize, cacheExpireMs,
maxRetryTimes);
- }
- }
-}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduRowDataLookupFunction.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduRowDataLookupFunction.java
index 4ae4121..8aeb6c8 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduRowDataLookupFunction.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/function/lookup/KuduRowDataLookupFunction.java
@@ -29,107 +29,74 @@ import
org.apache.flink.connector.kudu.connector.reader.KuduReaderIterator;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.functions.FunctionContext;
-import org.apache.flink.table.functions.TableFunction;
+import org.apache.flink.table.functions.LookupFunction;
+import org.apache.flink.util.CollectionUtil;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.commons.compress.utils.Lists;
-import org.apache.commons.lang3.ArrayUtils;
-import org.apache.kudu.shaded.com.google.common.cache.Cache;
-import org.apache.kudu.shaded.com.google.common.cache.CacheBuilder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
-import java.util.Arrays;
+import java.util.Collection;
+import java.util.Collections;
import java.util.List;
-import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
/** LookupFunction based on the RowData object type. */
-public class KuduRowDataLookupFunction extends TableFunction<RowData> {
+public class KuduRowDataLookupFunction extends LookupFunction {
private static final long serialVersionUID = 1L;
private static final Logger LOG =
LoggerFactory.getLogger(KuduRowDataLookupFunction.class);
private final KuduTableInfo tableInfo;
private final KuduReaderConfig kuduReaderConfig;
private final String[] keyNames;
- private final String[] projectedFields;
- private final long cacheMaxSize;
- private final long cacheExpireMs;
+ private final List<String> projectedFields;
private final int maxRetryTimes;
private final RowResultConverter<RowData> convertor;
- private transient Cache<RowData, List<RowData>> cache;
private transient KuduReader<RowData> kuduReader;
- private KuduRowDataLookupFunction(
+ public KuduRowDataLookupFunction(
String[] keyNames,
KuduTableInfo tableInfo,
KuduReaderConfig kuduReaderConfig,
- String[] projectedFields,
- KuduLookupOptions kuduLookupOptions) {
+ List<String> projectedFields) {
+ this(keyNames, tableInfo, kuduReaderConfig, projectedFields, 1);
+ }
+
+ public KuduRowDataLookupFunction(
+ String[] keyNames,
+ KuduTableInfo tableInfo,
+ KuduReaderConfig kuduReaderConfig,
+ List<String> projectedFields,
+ int maxRetryTimes) {
this.tableInfo = tableInfo;
- this.convertor = new RowResultRowDataConverter();
this.projectedFields = projectedFields;
this.keyNames = keyNames;
this.kuduReaderConfig = kuduReaderConfig;
- this.cacheMaxSize = kuduLookupOptions.getCacheMaxSize();
- this.cacheExpireMs = kuduLookupOptions.getCacheExpireMs();
- this.maxRetryTimes = kuduLookupOptions.getMaxRetryTimes();
- }
-
- public RowData buildCacheKey(Object... keys) {
- return GenericRowData.of(keys);
+ this.maxRetryTimes = maxRetryTimes;
+ convertor = new RowResultRowDataConverter();
}
- /**
- * invoke entry point of lookup function.
- *
- * @param keys join keys
- */
- public void eval(Object... keys) {
- if (keys.length != keyNames.length) {
- throw new RuntimeException("The join keys are of unequal lengths");
- }
- // cache key
- RowData keyRow = buildCacheKey(keys);
- if (this.cache != null) {
- List<RowData> cacheRows = this.cache.getIfPresent(keyRow);
- if (CollectionUtils.isNotEmpty(cacheRows)) {
- for (RowData cacheRow : cacheRows) {
- collect(cacheRow);
- }
- return;
- }
- }
-
+ @Override
+ public Collection<RowData> lookup(RowData keyRow) {
for (int retry = 1; retry <= maxRetryTimes; retry++) {
try {
- List<KuduFilterInfo> kuduFilterInfos =
buildKuduFilterInfo(keys);
+ List<KuduFilterInfo> kuduFilterInfos =
buildKuduFilterInfo((GenericRowData) keyRow);
this.kuduReader.setTableFilters(kuduFilterInfos);
KuduInputSplit[] inputSplits = kuduReader.createInputSplits(1);
ArrayList<RowData> rows = new ArrayList<>();
for (KuduInputSplit inputSplit : inputSplits) {
KuduReaderIterator<RowData> scanner =
kuduReader.scanner(inputSplit.getScanToken());
- // not use cache
- if (cache == null) {
- while (scanner.hasNext()) {
- collect(scanner.next());
- }
- } else {
- while (scanner.hasNext()) {
- RowData row = scanner.next();
- rows.add(row);
- collect(row);
- }
- rows.trimToSize();
+ while (scanner.hasNext()) {
+ RowData row = scanner.next();
+ rows.add(row);
}
}
- if (cache != null) {
- cache.put(keyRow, rows);
- }
- break;
+ rows.trimToSize();
+ return rows;
} catch (Exception e) {
LOG.error(String.format("Kudu scan error, retry times = %d",
retry), e);
if (retry >= maxRetryTimes) {
@@ -142,16 +109,8 @@ public class KuduRowDataLookupFunction extends
TableFunction<RowData> {
}
}
}
- }
- private List<KuduFilterInfo> buildKuduFilterInfo(Object... keyValS) {
- List<KuduFilterInfo> kuduFilterInfos = Lists.newArrayList();
- for (int i = 0; i < keyNames.length; i++) {
- KuduFilterInfo kuduFilterInfo =
-
KuduFilterInfo.Builder.create(keyNames[i]).equalTo(keyValS[i]).build();
- kuduFilterInfos.add(kuduFilterInfo);
- }
- return kuduFilterInfos;
+ return Collections.emptyList();
}
@Override
@@ -162,14 +121,7 @@ public class KuduRowDataLookupFunction extends
TableFunction<RowData> {
new KuduReader<>(this.tableInfo, this.kuduReaderConfig,
this.convertor);
// build kudu cache
this.kuduReader.setTableProjections(
- ArrayUtils.isNotEmpty(projectedFields) ?
Arrays.asList(projectedFields) : null);
- this.cache =
- this.cacheMaxSize == -1 || this.cacheExpireMs == -1
- ? null
- : CacheBuilder.newBuilder()
- .expireAfterWrite(this.cacheExpireMs,
TimeUnit.MILLISECONDS)
- .maximumSize(this.cacheMaxSize)
- .build();
+ CollectionUtil.isNullOrEmpty(projectedFields) ? null :
projectedFields);
} catch (Exception ioe) {
LOG.error("Exception while creating connection to Kudu.", ioe);
throw new RuntimeException("Cannot create connection to Kudu.",
ioe);
@@ -178,62 +130,24 @@ public class KuduRowDataLookupFunction extends
TableFunction<RowData> {
@Override
public void close() {
- if (null != this.kuduReader) {
+ if (kuduReader != null) {
try {
- this.kuduReader.close();
- if (cache != null) {
- this.cache.cleanUp();
- // help gc
- this.cache = null;
- }
- this.kuduReader = null;
+ kuduReader.close();
+ kuduReader = null;
} catch (IOException e) {
// ignore exception when close.
- LOG.warn("exception when close table", e);
+ LOG.warn("Failed to close Kudu table reader", e);
}
}
}
- /** Builder for KuduRowDataLookupFunction. */
- public static class Builder {
- private KuduTableInfo tableInfo;
- private KuduReaderConfig kuduReaderConfig;
- private String[] keyNames;
- private String[] projectedFields;
- private KuduLookupOptions kuduLookupOptions;
-
- public static Builder options() {
- return new Builder();
- }
-
- public Builder tableInfo(KuduTableInfo tableInfo) {
- this.tableInfo = tableInfo;
- return this;
- }
-
- public Builder kuduReaderConfig(KuduReaderConfig kuduReaderConfig) {
- this.kuduReaderConfig = kuduReaderConfig;
- return this;
- }
-
- public Builder keyNames(String[] keyNames) {
- this.keyNames = keyNames;
- return this;
- }
-
- public Builder projectedFields(String[] projectedFields) {
- this.projectedFields = projectedFields;
- return this;
- }
-
- public Builder kuduLookupOptions(KuduLookupOptions kuduLookupOptions) {
- this.kuduLookupOptions = kuduLookupOptions;
- return this;
- }
-
- public KuduRowDataLookupFunction build() {
- return new KuduRowDataLookupFunction(
- keyNames, tableInfo, kuduReaderConfig, projectedFields,
kuduLookupOptions);
- }
+ private List<KuduFilterInfo> buildKuduFilterInfo(GenericRowData keyRow) {
+ return IntStream.range(0, keyNames.length)
+ .mapToObj(
+ i ->
+ KuduFilterInfo.Builder.create(keyNames[i])
+ .equalTo(keyRow.getField(i))
+ .build())
+ .collect(Collectors.toList());
}
}
diff --git
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/utils/KuduTableUtils.java
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/utils/KuduTableUtils.java
index 63d721e..d8c259c 100644
---
a/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/utils/KuduTableUtils.java
+++
b/flink-connector-kudu/src/main/java/org/apache/flink/connector/kudu/table/utils/KuduTableUtils.java
@@ -22,8 +22,7 @@ import
org.apache.flink.connector.kudu.connector.ColumnSchemasFactory;
import org.apache.flink.connector.kudu.connector.CreateTableOptionsFactory;
import org.apache.flink.connector.kudu.connector.KuduFilterInfo;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
-import
org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory;
-import org.apache.flink.table.api.TableSchema;
+import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.expressions.CallExpression;
import org.apache.flink.table.expressions.Expression;
import org.apache.flink.table.expressions.FieldReferenceExpression;
@@ -32,8 +31,6 @@ import
org.apache.flink.table.functions.BuiltInFunctionDefinitions;
import org.apache.flink.table.functions.FunctionDefinition;
import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.logical.DecimalType;
-import org.apache.flink.table.types.logical.TimestampType;
-import org.apache.flink.table.utils.TableSchemaUtils;
import org.apache.kudu.ColumnSchema;
import org.apache.kudu.ColumnTypeAttributes;
@@ -44,7 +41,6 @@ import org.slf4j.LoggerFactory;
import javax.annotation.Nullable;
-import java.sql.Timestamp;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -53,42 +49,40 @@ import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_HASH_PARTITION_NUMS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_PRIMARY_KEY_COLS;
+import static
org.apache.flink.connector.kudu.table.KuduDynamicTableOptions.KUDU_REPLICAS;
+
/** Kudu table utilities. */
public class KuduTableUtils {
private static final Logger LOG =
LoggerFactory.getLogger(KuduTableUtils.class);
public static KuduTableInfo createTableInfo(
- String tableName, TableSchema schema, Map<String, String> props) {
+ String tableName, ResolvedSchema schema, Map<String, String>
props) {
// Since KUDU_HASH_COLS is a required property for table creation, we
use it to infer
// whether to create table
boolean createIfMissing =
-
props.containsKey(KuduDynamicTableSourceSinkFactory.KUDU_PRIMARY_KEY_COLS.key())
+ props.containsKey(KUDU_PRIMARY_KEY_COLS.key())
|| schema.getPrimaryKey().isPresent();
KuduTableInfo tableInfo = KuduTableInfo.forTable(tableName);
if (createIfMissing) {
-
List<Tuple2<String, DataType>> columns =
-
getSchemaWithSqlTimestamp(schema).getTableColumns().stream()
- .map(tc -> Tuple2.of(tc.getName(), tc.getType()))
+ schema.getColumns().stream()
+ .map(col -> Tuple2.of(col.getName(),
col.getDataType()))
.collect(Collectors.toList());
-
List<String> keyColumns = getPrimaryKeyColumns(props, schema);
+
ColumnSchemasFactory schemasFactory = () ->
toKuduConnectorColumns(columns, keyColumns);
int replicas =
- Optional.ofNullable(
- props.get(
-
KuduDynamicTableSourceSinkFactory.KUDU_REPLICAS.key()))
+ Optional.ofNullable(props.get(KUDU_REPLICAS.key()))
.map(Integer::parseInt)
.orElse(1);
// if hash partitions nums not exists,default 3;
int hashPartitionNums =
- Optional.ofNullable(
- props.get(
- KuduDynamicTableSourceSinkFactory
- .KUDU_HASH_PARTITION_NUMS
- .key()))
+
Optional.ofNullable(props.get(KUDU_HASH_PARTITION_NUMS.key()))
.map(Integer::parseInt)
.orElse(3);
CreateTableOptionsFactory optionsFactory =
@@ -100,7 +94,7 @@ public class KuduTableUtils {
} else {
LOG.debug(
"Property {} is missing, assuming the table is already
created.",
- KuduDynamicTableSourceSinkFactory.KUDU_HASH_COLS.key());
+ KUDU_HASH_COLS.key());
}
return tableInfo;
@@ -131,52 +125,29 @@ public class KuduTableUtils {
.collect(Collectors.toList());
}
- public static TableSchema kuduToFlinkSchema(Schema schema) {
- TableSchema.Builder builder = TableSchema.builder();
+ public static org.apache.flink.table.api.Schema kuduToFlinkSchema(Schema
schema) {
+ org.apache.flink.table.api.Schema.Builder builder =
+ org.apache.flink.table.api.Schema.newBuilder();
for (ColumnSchema column : schema.getColumns()) {
DataType flinkType =
KuduTypeUtils.toFlinkType(column.getType(),
column.getTypeAttributes())
.nullable();
- builder.field(column.getName(), flinkType);
+ builder.column(column.getName(), flinkType);
}
return builder.build();
}
public static List<String> getPrimaryKeyColumns(
- Map<String, String> tableProperties, TableSchema tableSchema) {
- return tableProperties.containsKey(
-
KuduDynamicTableSourceSinkFactory.KUDU_PRIMARY_KEY_COLS.key())
- ? Arrays.asList(
- tableProperties
-
.get(KuduDynamicTableSourceSinkFactory.KUDU_PRIMARY_KEY_COLS.key())
- .split(","))
- : tableSchema.getPrimaryKey().get().getColumns();
+ Map<String, String> tableProperties, ResolvedSchema schema) {
+ return tableProperties.containsKey(KUDU_PRIMARY_KEY_COLS.key())
+ ?
Arrays.asList(tableProperties.get(KUDU_PRIMARY_KEY_COLS.key()).split(","))
+ : schema.getPrimaryKey().get().getColumns();
}
public static List<String> getHashColumns(Map<String, String>
tableProperties) {
- return Arrays.asList(
- tableProperties
-
.get(KuduDynamicTableSourceSinkFactory.KUDU_HASH_COLS.key())
- .split(","));
- }
-
- public static TableSchema getSchemaWithSqlTimestamp(TableSchema schema) {
- TableSchema.Builder builder = new TableSchema.Builder();
- TableSchemaUtils.getPhysicalSchema(schema)
- .getTableColumns()
- .forEach(
- tableColumn -> {
- if (tableColumn.getType().getLogicalType()
instanceof TimestampType) {
- builder.field(
- tableColumn.getName(),
-
tableColumn.getType().bridgedTo(Timestamp.class));
- } else {
- builder.field(tableColumn.getName(),
tableColumn.getType());
- }
- });
- return builder.build();
+ return
Arrays.asList(tableProperties.get(KUDU_HASH_COLS.key()).split(","));
}
/** Converts Flink Expression to KuduFilterInfo. */
@@ -196,7 +167,7 @@ public class KuduTableUtils {
} else if (children.size() == 2
&&
!functionDefinition.equals(BuiltInFunctionDefinitions.OR)) {
return convertBinaryComparison(functionDefinition, children);
- } else if (children.size() > 0
+ } else if (!children.isEmpty()
&&
functionDefinition.equals(BuiltInFunctionDefinitions.OR)) {
return convertIsInExpression(children);
}
diff --git
a/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
b/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
index 48d5746..41d4de6 100644
---
a/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
+++
b/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
@@ -13,4 +13,5 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-org.apache.flink.connector.kudu.table.dynamic.KuduDynamicTableSourceSinkFactory
+org.apache.flink.connector.kudu.table.KuduDynamicTableFactory
+org.apache.flink.connector.kudu.table.catalog.KuduCatalogFactory
diff --git
a/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
b/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
deleted file mode 100644
index ca6fe21..0000000
---
a/flink-connector-kudu/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
+++ /dev/null
@@ -1,17 +0,0 @@
-# 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.
-
-org.apache.flink.connector.kudu.table.KuduTableFactory
-org.apache.flink.connector.kudu.table.dynamic.catalog.KuduCatalogFactory
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/connector/KuduTestBase.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/connector/KuduTestBase.java
index 0b98fa4..c9e3755 100644
---
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/connector/KuduTestBase.java
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/connector/KuduTestBase.java
@@ -29,7 +29,8 @@ import
org.apache.flink.connector.kudu.connector.writer.KuduWriter;
import org.apache.flink.connector.kudu.connector.writer.KuduWriterConfig;
import org.apache.flink.connector.kudu.connector.writer.RowOperationMapper;
import org.apache.flink.table.api.DataTypes;
-import org.apache.flink.table.api.TableSchema;
+import org.apache.flink.table.catalog.Column;
+import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.StringData;
@@ -199,14 +200,13 @@ public class KuduTestBase {
.collect(Collectors.toList());
}
- public static TableSchema booksTableSchema() {
- return TableSchema.builder()
- .field("id", DataTypes.INT())
- .field("title", DataTypes.STRING())
- .field("author", DataTypes.STRING())
- .field("price", DataTypes.DOUBLE())
- .field("quantity", DataTypes.INT())
- .build();
+ public static ResolvedSchema booksTableSchema() {
+ return ResolvedSchema.of(
+ Column.physical("id", DataTypes.INT()),
+ Column.physical("title", DataTypes.STRING()),
+ Column.physical("author", DataTypes.STRING()),
+ Column.physical("price", DataTypes.DOUBLE()),
+ Column.physical("quantity", DataTypes.INT()));
}
public static List<RowData> booksRowData() {
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSinkTest.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSinkTest.java
similarity index 93%
rename from
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSinkTest.java
rename to
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSinkTest.java
index 5b417c8..3e68a12 100644
---
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSinkTest.java
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSinkTest.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic;
+package org.apache.flink.connector.kudu.table;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import org.apache.flink.connector.kudu.connector.KuduTestBase;
@@ -31,7 +31,7 @@ import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertNotNull;
-/** Unit Tests for {@link KuduDynamicTableSink}. */
+/** Unit Tests for {@link
org.apache.flink.connector.kudu.table.KuduDynamicTableSink}. */
public class KuduDynamicSinkTest extends KuduTestBase {
public static final String INPUT_TABLE = "books";
public static StreamExecutionEnvironment env;
@@ -52,7 +52,7 @@ public class KuduDynamicSinkTest extends KuduTestBase {
}
@Test
- public void testKuduSink() throws Exception {
+ public void testKuduSink() {
String createSql =
"CREATE TABLE "
+ INPUT_TABLE
@@ -74,7 +74,7 @@ public class KuduDynamicSinkTest extends KuduTestBase {
+ "','kudu.flush-interval'='1000"
+ "','kudu.operation-timeout'='500"
+ "','kudu.ignore-not-found'='true"
- + "','kudu.ignore-not-found'='true'"
+ + "','kudu.ignore-duplicate'='true'"
+ ")";
tEnv.executeSql(createSql);
tEnv.executeSql(
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSourceTest.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSourceTest.java
similarity index 98%
rename from
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSourceTest.java
rename to
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSourceTest.java
index 967d61c..ad0be6d 100644
---
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduDynamicSourceTest.java
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicSourceTest.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic;
+package org.apache.flink.connector.kudu.table;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import org.apache.flink.connector.kudu.connector.KuduTestBase;
@@ -37,7 +37,7 @@ import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
-/** Unit Tests for {@link KuduDynamicTableSource}. */
+/** Unit Tests for {@link
org.apache.flink.connector.kudu.table.KuduDynamicTableSource}. */
public class KuduDynamicSourceTest extends KuduTestBase {
public static final String INPUT_TABLE = "books";
public static StreamExecutionEnvironment env;
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactoryTest.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactoryTest.java
new file mode 100644
index 0000000..46abe1a
--- /dev/null
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduDynamicTableFactoryTest.java
@@ -0,0 +1,233 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.connector.kudu.connector.KuduTableInfo;
+import org.apache.flink.connector.kudu.connector.KuduTestBase;
+import org.apache.flink.connector.kudu.connector.writer.KuduWriterConfig;
+import org.apache.flink.core.execution.JobClient;
+import org.apache.flink.runtime.client.JobExecutionException;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.api.ValidationException;
+import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
+import org.apache.flink.table.catalog.CatalogTable;
+import org.apache.flink.table.catalog.Column;
+import org.apache.flink.table.catalog.ObjectIdentifier;
+import org.apache.flink.table.catalog.ResolvedCatalogTable;
+import org.apache.flink.table.catalog.ResolvedSchema;
+import org.apache.flink.table.connector.sink.DynamicTableSink;
+import org.apache.flink.table.factories.FactoryUtil;
+
+import org.apache.kudu.Type;
+import org.apache.kudu.client.KuduScanner;
+import org.apache.kudu.client.KuduTable;
+import org.apache.kudu.client.RowResult;
+import org.apache.kudu.shaded.com.google.common.collect.Lists;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.sql.Timestamp;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.fail;
+
+/** Tests for dynamic table factory. */
+public class KuduDynamicTableFactoryTest extends KuduTestBase {
+
+ private StreamTableEnvironment tableEnv;
+ private String kuduMasters;
+
+ @BeforeEach
+ public void init() {
+ StreamExecutionEnvironment env =
+
StreamExecutionEnvironment.getExecutionEnvironment().setParallelism(1);
+ tableEnv = KuduTableTestUtils.createTableEnvInStreamingMode(env);
+ kuduMasters = getMasterAddress();
+ }
+
+ @Test
+ public void testMissingMasters() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable11 (`first` STRING, `second` INT) "
+ + "WITH ('connector'='kudu',
'kudu.table'='TestTable11')");
+ assertThrows(
+ ValidationException.class,
+ () -> tableEnv.executeSql("INSERT INTO TestTable11 values
('f', 1)"));
+ }
+
+ @Test
+ public void testNonExistingTable() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable11 (`first` STRING, `second` INT) "
+ + "WITH ('connector'='kudu',
'kudu.table'='TestTable11', 'kudu.masters'='"
+ + kuduMasters
+ + "')");
+ JobClient jobClient =
+ tableEnv.executeSql("INSERT INTO TestTable11 values ('f',
1)").getJobClient().get();
+ try {
+ jobClient.getJobExecutionResult().get();
+ fail();
+ } catch (ExecutionException ee) {
+ assertTrue(ee.getCause() instanceof JobExecutionException);
+ }
+ }
+
+ @Test
+ public void testCreateTable() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable11 (`first` STRING, `second` STRING) "
+ + "WITH ('connector'='kudu',
'kudu.table'='TestTable11', 'kudu.masters'='"
+ + kuduMasters
+ + "', "
+ + "'kudu.hash-columns'='first',
'kudu.primary-key-columns'='first')");
+
+ tableEnv.executeSql("INSERT INTO TestTable11 values ('f', 's')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ validateSingleKey("TestTable11");
+ }
+
+ @Test
+ public void testTimestamp() throws Exception {
+ // Timestamp should be bridged to sql.Timestamp
+ // Test it when creating the table...
+ tableEnv.executeSql(
+ "CREATE TABLE TestTableTs (`first` STRING, `second`
TIMESTAMP(3)) "
+ + "WITH ('connector'='kudu', 'kudu.masters'='"
+ + kuduMasters
+ + "', "
+ + "'kudu.hash-columns'='first',
'kudu.primary-key-columns'='first')");
+ tableEnv.executeSql(
+ "INSERT INTO TestTableTs values ('f', TIMESTAMP
'2020-01-01 12:12:12.123456')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ tableEnv.executeSql("INSERT INTO TestTableTs values ('s', TIMESTAMP
'2020-02-02 23:23:23')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ KuduTable kuduTable = getClient().openTable("TestTableTs");
+ assertEquals(Type.UNIXTIME_MICROS,
kuduTable.getSchema().getColumn("second").getType());
+
+ KuduScanner scanner = getClient().newScannerBuilder(kuduTable).build();
+ HashSet<Timestamp> results = new HashSet<>();
+ scanner.forEach(sc -> results.add(sc.getTimestamp("second")));
+
+ assertEquals(2, results.size());
+ List<Timestamp> expected =
+ Lists.newArrayList(
+ Timestamp.valueOf("2020-01-01 12:12:12.123"),
+ Timestamp.valueOf("2020-02-02 23:23:23"));
+ assertEquals(new HashSet<>(expected), results);
+ }
+
+ @Test
+ public void testExistingTable() throws Exception {
+ // Creating a table
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable12 (`first` STRING, `second` STRING) "
+ + "WITH ('connector'='kudu',
'kudu.table'='TestTable12', 'kudu.masters'='"
+ + kuduMasters
+ + "', "
+ + "'kudu.hash-columns'='first',
'kudu.primary-key-columns'='first')");
+
+ tableEnv.executeSql("INSERT INTO TestTable12 values ('f', 's')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ // Then another one in SQL that refers to the previously created one
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable12b (`first` STRING, `second` STRING) "
+ + "WITH ('connector'='kudu',
'kudu.table'='TestTable12', 'kudu.masters'='"
+ + kuduMasters
+ + "')");
+ tableEnv.executeSql("INSERT INTO TestTable12b values ('f2','s2')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ // Validate that both insertions were into the same table
+ KuduTable kuduTable = getClient().openTable("TestTable12");
+ KuduScanner scanner = getClient().newScannerBuilder(kuduTable).build();
+ List<RowResult> rows = new ArrayList<>();
+ scanner.forEach(rows::add);
+
+ assertEquals(2, rows.size());
+ assertEquals("f", rows.get(0).getString("first"));
+ assertEquals("s", rows.get(0).getString("second"));
+ assertEquals("f2", rows.get(1).getString("first"));
+ assertEquals("s2", rows.get(1).getString("second"));
+ }
+
+ @Test
+ public void testTableSink() {
+ final ResolvedSchema schema =
+ ResolvedSchema.of(
+ Column.physical("first", DataTypes.STRING()),
+ Column.physical("second", DataTypes.STRING()));
+ final Map<String, String> properties = new HashMap<>();
+ properties.put("connector", "kudu");
+ properties.put("kudu.masters", kuduMasters);
+ properties.put("kudu.table", "TestTable12");
+ properties.put("kudu.ignore-not-found", "true");
+ properties.put("kudu.ignore-duplicate", "true");
+ properties.put("kudu.flush-interval", "10000");
+ properties.put("kudu.max-buffer-size", "10000");
+
+ KuduWriterConfig.Builder builder =
+ KuduWriterConfig.Builder.setMasters(kuduMasters)
+ .setFlushInterval(10000)
+ .setMaxBufferSize(10000)
+ .setIgnoreDuplicate(true)
+ .setIgnoreNotFound(true);
+ KuduTableInfo kuduTableInfo = KuduTableInfo.forTable("TestTable12");
+ KuduDynamicTableSink expected = new KuduDynamicTableSink(builder,
kuduTableInfo, schema);
+ final DynamicTableSink actualSink =
+ FactoryUtil.createDynamicTableSink(
+ null,
+ ObjectIdentifier.of("kudu", "default", "TestTable12"),
+ new
ResolvedCatalogTable(CatalogTable.fromProperties(properties), schema),
+ properties,
+ new Configuration(),
+ Thread.currentThread().getContextClassLoader(),
+ false);
+
+ assertEquals(expected, actualSink);
+ }
+}
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduRowDataLookupFunctionTest.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduRowDataLookupFunctionTest.java
similarity index 62%
rename from
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduRowDataLookupFunctionTest.java
rename to
flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduRowDataLookupFunctionTest.java
index bc728ec..c93391c 100644
---
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/dynamic/KuduRowDataLookupFunctionTest.java
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduRowDataLookupFunctionTest.java
@@ -15,12 +15,11 @@
* limitations under the License.
*/
-package org.apache.flink.connector.kudu.table.dynamic;
+package org.apache.flink.connector.kudu.table;
import org.apache.flink.connector.kudu.connector.KuduTableInfo;
import org.apache.flink.connector.kudu.connector.KuduTestBase;
import org.apache.flink.connector.kudu.connector.reader.KuduReaderConfig;
-import org.apache.flink.connector.kudu.table.function.lookup.KuduLookupOptions;
import
org.apache.flink.connector.kudu.table.function.lookup.KuduRowDataLookupFunction;
import org.apache.flink.table.data.RowData;
import org.apache.flink.util.Collector;
@@ -30,6 +29,7 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;
@@ -53,41 +53,8 @@ public class KuduRowDataLookupFunctionTest extends
KuduTestBase {
}
@Test
- public void testEval() throws Exception {
- KuduLookupOptions lookupOptions = KuduLookupOptions.builder().build();
-
- KuduRowDataLookupFunction lookupFunction =
- buildRowDataLookupFunction(lookupOptions, new String[] {"id"});
-
- ListOutputCollector collector = new ListOutputCollector();
- lookupFunction.setCollector(collector);
-
- lookupFunction.open(null);
-
- lookupFunction.eval(1001);
-
- lookupFunction.eval(1002);
-
- lookupFunction.eval(1003);
-
- List<String> result =
- new ArrayList<>(collector.getOutputs())
-
.stream().map(RowData::toString).sorted().collect(Collectors.toList());
-
- assertNotNull(result);
- }
-
- @Test
- public void testCacheEval() throws Exception {
- KuduLookupOptions lookupOptions =
- KuduLookupOptions.builder()
- .withCacheMaxSize(1024)
- .withMaxRetryTimes(3)
- .withCacheExpireMs(10)
- .build();
-
- KuduRowDataLookupFunction lookupFunction =
- buildRowDataLookupFunction(lookupOptions, new String[] {"id"});
+ public void testLookup() throws Exception {
+ KuduRowDataLookupFunction lookupFunction =
buildRowDataLookupFunction(new String[] {"id"});
ListOutputCollector collector = new ListOutputCollector();
lookupFunction.setCollector(collector);
@@ -107,21 +74,14 @@ public class KuduRowDataLookupFunctionTest extends
KuduTestBase {
assertNotNull(result);
}
- private KuduRowDataLookupFunction buildRowDataLookupFunction(
- KuduLookupOptions lookupOptions, String[] keyNames) {
+ private KuduRowDataLookupFunction buildRowDataLookupFunction(String[]
keyNames) {
KuduReaderConfig config =
KuduReaderConfig.Builder.setMasters(getMasterAddress()).setRowLimit(10).build();
- return new KuduRowDataLookupFunction.Builder()
- .kuduReaderConfig(config)
- .kuduLookupOptions(lookupOptions)
- .keyNames(keyNames)
- .projectedFields(getFieldNames())
- .tableInfo(tableInfo)
- .build();
- }
-
- private String[] getFieldNames() {
- return new String[] {"id", "title", "author", "price", "quantity"};
+ return new KuduRowDataLookupFunction(
+ keyNames,
+ tableInfo,
+ config,
+ Arrays.asList("id", "title", "author", "price", "quantity"));
}
private static final class ListOutputCollector implements
Collector<RowData> {
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableSourceITCase.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableSourceITCase.java
new file mode 100644
index 0000000..87c8828
--- /dev/null
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableSourceITCase.java
@@ -0,0 +1,87 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.connector.kudu.connector.KuduTableInfo;
+import org.apache.flink.connector.kudu.connector.KuduTestBase;
+import org.apache.flink.connector.kudu.table.catalog.KuduCatalog;
+import org.apache.flink.table.api.TableEnvironment;
+import org.apache.flink.types.Row;
+import org.apache.flink.util.CloseableIterator;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/** Integration tests for {@link KuduDynamicTableSource}. */
+public class KuduTableSourceITCase extends KuduTestBase {
+ private TableEnvironment tableEnv;
+ private KuduCatalog catalog;
+
+ private KuduTableInfo tableInfo = null;
+
+ @BeforeEach
+ void init() {
+ tableInfo = booksTableInfo("books", true);
+ setUpDatabase(tableInfo);
+ tableEnv = KuduTableTestUtils.createTableEnvInBatchMode();
+ catalog = new KuduCatalog(getMasterAddress());
+ tableEnv.registerCatalog("kudu", catalog);
+ tableEnv.useCatalog("kudu");
+ }
+
+ @AfterEach
+ void cleanup() {
+ if (tableInfo != null) {
+ cleanDatabase(tableInfo);
+ tableInfo = null;
+ }
+ }
+
+ @Test
+ void testFullBatchScan() throws Exception {
+ CloseableIterator<Row> it =
+ tableEnv.executeSql("select * from books order by
id").collect();
+ List<Row> results = new ArrayList<>();
+ it.forEachRemaining(results::add);
+ assertEquals(5, results.size());
+ assertEquals(
+ "+I[1001, Java for dummies, Tan Ah Teck, 11.11, 11]",
results.get(0).toString());
+ tableEnv.executeSql("DROP TABLE books");
+ }
+
+ @Test
+ void testScanWithProjectionAndFilter() throws Exception {
+ // (price > 30 and price < 40)
+ CloseableIterator<Row> it =
+ tableEnv.executeSql(
+ "SELECT title FROM books WHERE id IN (1003,
1004) and "
+ + "quantity < 40")
+ .collect();
+ List<Row> results = new ArrayList<>();
+ it.forEachRemaining(results::add);
+ assertEquals(1, results.size());
+ assertEquals("+I[More Java for more dummies]",
results.get(0).toString());
+ tableEnv.executeSql("DROP TABLE books");
+ }
+}
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableTestUtils.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableTestUtils.java
new file mode 100644
index 0000000..737ac96
--- /dev/null
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/KuduTableTestUtils.java
@@ -0,0 +1,44 @@
+/*
+ * 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.flink.connector.kudu.table;
+
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.api.EnvironmentSettings;
+import org.apache.flink.table.api.TableEnvironment;
+import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
+
+import static
org.apache.flink.table.api.config.ExecutionConfigOptions.TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM;
+
+/** Table API test utilities. */
+public class KuduTableTestUtils {
+
+ public static StreamTableEnvironment createTableEnvInStreamingMode(
+ StreamExecutionEnvironment env) {
+ EnvironmentSettings settings =
EnvironmentSettings.newInstance().inStreamingMode().build();
+ StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env,
settings);
+
tableEnv.getConfig().getConfiguration().set(TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM,
1);
+ return tableEnv;
+ }
+
+ public static TableEnvironment createTableEnvInBatchMode() {
+ EnvironmentSettings settings =
EnvironmentSettings.newInstance().inBatchMode().build();
+ TableEnvironment tableEnv = TableEnvironment.create(settings);
+
tableEnv.getConfig().getConfiguration().set(TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM,
1);
+ return tableEnv;
+ }
+}
diff --git
a/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogTest.java
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogTest.java
new file mode 100644
index 0000000..d3a5f91
--- /dev/null
+++
b/flink-connector-kudu/src/test/java/org/apache/flink/connector/kudu/table/catalog/KuduCatalogTest.java
@@ -0,0 +1,292 @@
+/*
+ * 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.flink.connector.kudu.table.catalog;
+
+import org.apache.flink.api.common.typeinfo.Types;
+import org.apache.flink.api.java.tuple.Tuple2;
+import org.apache.flink.connector.kudu.connector.KuduTestBase;
+import org.apache.flink.connector.kudu.table.KuduTableTestUtils;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.SinkFunction;
+import org.apache.flink.table.api.Table;
+import org.apache.flink.table.api.TableException;
+import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
+import org.apache.flink.types.Row;
+
+import org.apache.kudu.Schema;
+import org.apache.kudu.Type;
+import org.apache.kudu.client.KuduScanner;
+import org.apache.kudu.client.KuduTable;
+import org.apache.kudu.client.RowResult;
+import org.apache.kudu.shaded.com.google.common.collect.Lists;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.nio.ByteBuffer;
+import java.sql.Timestamp;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/** Tests for {@link KuduCatalog}. */
+public class KuduCatalogTest extends KuduTestBase {
+
+ private KuduCatalog catalog;
+ private StreamTableEnvironment tableEnv;
+
+ @BeforeEach
+ public void init() {
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
+ catalog = new KuduCatalog(getMasterAddress());
+ tableEnv = KuduTableTestUtils.createTableEnvInStreamingMode(env);
+ tableEnv.registerCatalog("kudu", catalog);
+ tableEnv.useCatalog("kudu");
+ }
+
+ @Test
+ public void testCreateAlterDrop() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable1 (`first` STRING, `second` String)
WITH ('kudu.hash-columns' = 'first', 'kudu.primary-key-columns' = 'first')");
+ tableEnv.executeSql("INSERT INTO TestTable1 VALUES ('f', 's')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ // Add this once Primary key support has been enabled
+ // tableEnv.sqlUpdate("CREATE TABLE TestTable2 (`first` STRING,
`second` String, PRIMARY
+ // KEY(`first`)) WITH ('kudu.hash-columns' = 'first')");
+ // tableEnv.sqlUpdate("INSERT INTO TestTable2 VALUES ('f', 's')");
+
+ validateSingleKey("TestTable1");
+ // validateSingleKey("TestTable2");
+
+ tableEnv.executeSql("ALTER TABLE TestTable1 RENAME TO TestTable1R");
+ validateSingleKey("TestTable1R");
+
+ tableEnv.executeSql("DROP TABLE TestTable1R");
+ assertFalse(getClient().tableExists("TestTable1R"));
+ }
+
+ @Test
+ public void testCreateAndInsertMultiKey() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable3 (`first` STRING, `second` INT, third
STRING) WITH ('kudu.hash-columns' = 'first,second', 'kudu.primary-key-columns'
= 'first,second')");
+ tableEnv.executeSql("INSERT INTO TestTable3 VALUES ('f', 2, 't')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ validateMultiKey("TestTable3");
+ }
+
+ @Test
+ public void testSourceProjection() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable5 (`second` String, `first` STRING,
`third` String) WITH ('kudu.hash-columns' = 'second',
'kudu.primary-key-columns' = 'second')");
+ tableEnv.executeSql("INSERT INTO TestTable5 VALUES ('s', 'f', 't')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable6 (`first` STRING, `second` String)
WITH ('kudu.hash-columns' = 'first', 'kudu.primary-key-columns' = 'first')");
+ tableEnv.executeSql("INSERT INTO TestTable6 (SELECT `first`, `second`
FROM TestTable5)")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ validateSingleKey("TestTable6");
+ }
+
+ @Test
+ public void testEmptyProjection() throws Exception {
+ CollectionSink.output.clear();
+ tableEnv.executeSql(
+ "CREATE TABLE TestTableEP (`first` STRING, `second` STRING)
WITH ('kudu.hash-columns' = 'first', 'kudu.primary-key-columns' = 'first')");
+ tableEnv.executeSql("INSERT INTO TestTableEP VALUES ('f','s')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+ tableEnv.executeSql("INSERT INTO TestTableEP VALUES ('f2','s2')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ Table result = tableEnv.sqlQuery("SELECT COUNT(*) FROM TestTableEP");
+
+ DataStream<Tuple2<Boolean, Row>> resultDataStream =
+ tableEnv.toRetractStream(result, Types.ROW(Types.LONG));
+
+ resultDataStream
+ .map(t -> Tuple2.of(t.f0, t.f1.getField(0)))
+ .returns(Types.TUPLE(Types.BOOLEAN, Types.LONG))
+ .addSink(new CollectionSink<>())
+ .setParallelism(1);
+
+ resultDataStream.getExecutionEnvironment().execute();
+
+ List<Tuple2<Boolean, Long>> expected =
+ Lists.newArrayList(Tuple2.of(true, 1L), Tuple2.of(false, 1L),
Tuple2.of(true, 2L));
+
+ assertEquals(new HashSet<>(expected), new
HashSet<>(CollectionSink.output));
+ CollectionSink.output.clear();
+ }
+
+ @Test
+ public void testTimestamp() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTableTsC (`first` STRING, `second`
TIMESTAMP(3)) "
+ + "WITH ('kudu.hash-columns'='first',
'kudu.primary-key-columns'='first')");
+ tableEnv.executeSql(
+ "INSERT INTO TestTableTsC values ('f', TIMESTAMP
'2020-01-01 12:12:12.123456')")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ KuduTable kuduTable = getClient().openTable("TestTableTsC");
+ assertEquals(Type.UNIXTIME_MICROS,
kuduTable.getSchema().getColumn("second").getType());
+
+ KuduScanner scanner = getClient().newScannerBuilder(kuduTable).build();
+ List<RowResult> rows = new ArrayList<>();
+ scanner.forEach(rows::add);
+
+ assertEquals(1, rows.size());
+ assertEquals("f", rows.get(0).getString(0));
+ assertEquals(Timestamp.valueOf("2020-01-01 12:12:12.123"),
rows.get(0).getTimestamp(1));
+ }
+
+ @Test
+ public void testDatatypes() throws Exception {
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable8 (`first` STRING, `second` BOOLEAN,
`third` BYTES,"
+ + "`fourth` TINYINT, `fifth` SMALLINT, `sixth` INT,
`seventh` BIGINT, `eighth` FLOAT, `ninth` DOUBLE, "
+ + "`tenth` TIMESTAMP)"
+ + "WITH ('kudu.hash-columns' = 'first',
'kudu.primary-key-columns' = 'first')");
+
+ tableEnv.executeSql(
+ "INSERT INTO TestTable8 values ('f', false,
cast('bbbb' as BYTES), cast(12 as TINYINT),"
+ + "cast(34 as SMALLINT), 56, cast(78 as
BIGINT), cast(3.14 as FLOAT), cast(1.2345 as DOUBLE),"
+ + "TIMESTAMP '2020-04-15 12:34:56.123') ")
+ .getJobClient()
+ .get()
+ .getJobExecutionResult()
+ .get(1, TimeUnit.MINUTES);
+
+ validateManyTypes("TestTable8");
+ }
+
+ @Test
+ public void testMissingPropertiesCatalog() throws Exception {
+ assertThrows(
+ TableException.class,
+ () ->
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable9a (`first` STRING,
`second` String) "
+ + "WITH ('kudu.primary-key-columns' =
'second')"));
+ assertThrows(
+ TableException.class,
+ () ->
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable9b (`first` STRING,
`second` String) "
+ + "WITH ('kudu.hash-columns' =
'first')"));
+ assertThrows(
+ TableException.class,
+ () ->
+ tableEnv.executeSql(
+ "CREATE TABLE TestTable9b (`first` STRING,
`second` String) "
+ + "WITH ('kudu.primary-key-columns' =
'second', 'kudu.hash-columns' = 'first')"));
+ }
+
+ private void validateManyTypes(String tableName) throws Exception {
+ KuduTable kuduTable = getClient().openTable(tableName);
+ Schema schema = kuduTable.getSchema();
+
+ assertEquals(Type.STRING, schema.getColumn("first").getType());
+ assertEquals(Type.BOOL, schema.getColumn("second").getType());
+ assertEquals(Type.BINARY, schema.getColumn("third").getType());
+ assertEquals(Type.INT8, schema.getColumn("fourth").getType());
+ assertEquals(Type.INT16, schema.getColumn("fifth").getType());
+ assertEquals(Type.INT32, schema.getColumn("sixth").getType());
+ assertEquals(Type.INT64, schema.getColumn("seventh").getType());
+ assertEquals(Type.FLOAT, schema.getColumn("eighth").getType());
+ assertEquals(Type.DOUBLE, schema.getColumn("ninth").getType());
+ assertEquals(Type.UNIXTIME_MICROS,
schema.getColumn("tenth").getType());
+
+ KuduScanner scanner = getClient().newScannerBuilder(kuduTable).build();
+ List<RowResult> rows = new ArrayList<>();
+ scanner.forEach(rows::add);
+
+ assertEquals(1, rows.size());
+ assertEquals("f", rows.get(0).getString(0));
+ assertEquals(false, rows.get(0).getBoolean(1));
+ assertEquals(ByteBuffer.wrap("bbbb".getBytes()),
rows.get(0).getBinary(2));
+ assertEquals(12, rows.get(0).getByte(3));
+ assertEquals(34, rows.get(0).getShort(4));
+ assertEquals(56, rows.get(0).getInt(5));
+ assertEquals(78, rows.get(0).getLong(6));
+ assertEquals(3.14, rows.get(0).getFloat(7), 0.01);
+ assertEquals(1.2345, rows.get(0).getDouble(8), 0.0001);
+ assertEquals(Timestamp.valueOf("2020-04-15 12:34:56.123"),
rows.get(0).getTimestamp(9));
+ }
+
+ private void validateMultiKey(String tableName) throws Exception {
+ KuduTable kuduTable = getClient().openTable(tableName);
+ Schema schema = kuduTable.getSchema();
+
+ assertEquals(2, schema.getPrimaryKeyColumnCount());
+ assertEquals(3, schema.getColumnCount());
+
+ assertTrue(schema.getColumn("first").isKey());
+ assertTrue(schema.getColumn("second").isKey());
+
+ assertFalse(schema.getColumn("third").isKey());
+
+ KuduScanner scanner = getClient().newScannerBuilder(kuduTable).build();
+ List<RowResult> rows = new ArrayList<>();
+ scanner.forEach(rows::add);
+
+ assertEquals(1, rows.size());
+ assertEquals("f", rows.get(0).getString("first"));
+ assertEquals(2, rows.get(0).getInt("second"));
+ assertEquals("t", rows.get(0).getString("third"));
+ }
+
+ private static class CollectionSink<T> implements SinkFunction<T> {
+
+ public static List<Object> output = Collections.synchronizedList(new
ArrayList<>());
+
+ public void invoke(T value, SinkFunction.Context context) {
+ output.add(value);
+ }
+ }
+}