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

Reply via email to