From a021b12b7ac5f12d5855ec4e7b88b8bd87c8e836 Mon Sep 17 00:00:00 2001 From: zhangjun Date: Thu, 6 May 2021 18:00:03 +0800 Subject: [PATCH 1/2] upgrade to flink 1.13 --- .../iceberg/flink/FlinkCatalogFactory.java | 42 +++++++--- .../flink/FlinkCatalogFactoryOptions.java | 84 +++++++++++++++++++ ... org.apache.flink.table.factories.Factory} | 2 +- versions.props | 2 +- 4 files changed, 115 insertions(+), 15 deletions(-) create mode 100644 flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java rename flink/src/main/resources/META-INF/services/{org.apache.flink.table.factories.TableFactory => org.apache.flink.table.factories.Factory} (94%) diff --git a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java index 406c9793021b..ff93bee78aaa 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java +++ b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java @@ -22,22 +22,22 @@ import java.net.URL; import java.nio.file.Files; import java.nio.file.Paths; -import java.util.List; +import java.util.HashSet; import java.util.Locale; import java.util.Map; +import java.util.Set; +import org.apache.flink.configuration.ConfigOption; import org.apache.flink.configuration.GlobalConfiguration; import org.apache.flink.runtime.util.HadoopUtils; import org.apache.flink.table.catalog.Catalog; -import org.apache.flink.table.descriptors.CatalogDescriptorValidator; import org.apache.flink.table.factories.CatalogFactory; +import org.apache.flink.table.factories.FactoryUtil; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.iceberg.CatalogProperties; import org.apache.iceberg.catalog.Namespace; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.base.Strings; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; -import org.apache.iceberg.relocated.com.google.common.collect.Maps; /** * A Flink Catalog factory implementation that creates {@link FlinkCatalog}. @@ -100,21 +100,37 @@ protected CatalogLoader createCatalogLoader(String name, Map pro } @Override - public Map requiredContext() { - Map context = Maps.newHashMap(); - context.put(CatalogDescriptorValidator.CATALOG_TYPE, "iceberg"); - context.put(CatalogDescriptorValidator.CATALOG_PROPERTY_VERSION, "1"); - return context; + public String factoryIdentifier() { + return FlinkCatalogFactoryOptions.IDENTIFIER; } @Override - public List supportedProperties() { - return ImmutableList.of("*"); + public Set> requiredOptions() { + final Set> options = new HashSet<>(); + options.add(FlinkCatalogFactoryOptions.CATALOG_TYPE); + return options; } @Override - public Catalog createCatalog(String name, Map properties) { - return createCatalog(name, properties, clusterHadoopConf()); + public Set> optionalOptions() { + final Set> options = new HashSet<>(); + options.add(FlinkCatalogFactoryOptions.PROPERTY_VERSION); + options.add(FlinkCatalogFactoryOptions.URI); + options.add(FlinkCatalogFactoryOptions.WAREHOUSE); + options.add(FlinkCatalogFactoryOptions.CLIENTS); + options.add(FlinkCatalogFactoryOptions.BASE_NAMESPACES); + options.add(FlinkCatalogFactoryOptions.HIVE_CONF_DIF); + options.add(FlinkCatalogFactoryOptions.CACHE_ENABLED); + return options; + } + + @Override + public Catalog createCatalog(Context context) { + final FactoryUtil.CatalogFactoryHelper helper = + FactoryUtil.createCatalogFactoryHelper(this, context); + helper.validate(); + + return createCatalog(context.getName(), context.getOptions(), clusterHadoopConf()); } protected Catalog createCatalog(String name, Map properties, Configuration hadoopConf) { diff --git a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java new file mode 100644 index 000000000000..3ba913d22453 --- /dev/null +++ b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java @@ -0,0 +1,84 @@ +/* + * 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.iceberg.flink; + +import org.apache.flink.configuration.ConfigOption; +import org.apache.flink.configuration.ConfigOptions; + +public class FlinkCatalogFactoryOptions { + private FlinkCatalogFactoryOptions() { + } + + public static final String IDENTIFIER = "iceberg"; + + public static final ConfigOption CATALOG_TYPE = + ConfigOptions.key("catalog-type") + .stringType() + .noDefaultValue() + .withDescription("Iceberg catalog type, 'hive' or 'hadoop'"); + + public static final ConfigOption PROPERTY_VERSION = + ConfigOptions.key("property-version") + .stringType() + .defaultValue("1") + .withDescription( + "Version number to describe the property version. This property can be used for backwards " + + "compatibility in case the property format changes. The current property version is `1`. (Optional)"); + + public static final ConfigOption URI = + ConfigOptions.key("uri") + .stringType() + .noDefaultValue() + .withDescription("The Hive Metastore URI (Hive catalog only)"); + + public static final ConfigOption CLIENTS = + ConfigOptions.key("clients") + .stringType() + .noDefaultValue() + .withDescription("The Hive Client Pool Size (Hive catalog only)"); + + public static final ConfigOption HIVE_CONF_DIF = + ConfigOptions.key("hive-conf-dir") + .stringType() + .noDefaultValue() + .withDescription( + "Path to a directory containing a hive-site.xml configuration file which will be used to provide " + + "custom Hive configuration values. The value of hive.metastore.warehouse.dir from" + + " /hive-site.xml (or hive configure file from classpath) will be overwrote with " + + "the warehouse value if setting both hive-conf-dir and warehouse when creating iceberg catalog."); + + public static final ConfigOption BASE_NAMESPACES = + ConfigOptions.key("base-namespace") + .stringType() + .noDefaultValue() + .withDescription("A base namespace as the prefix for all databases (Hadoop catalog only)"); + + public static final ConfigOption WAREHOUSE = + ConfigOptions.key("warehouse") + .stringType() + .noDefaultValue() + .withDescription("The warehouse path (Hadoop catalog only)"); + + public static final ConfigOption CACHE_ENABLED = + ConfigOptions.key("cache-enabled") + .stringType() + .defaultValue("true") + .withDescription("Whether to cache the catalog in FlinkCatalog."); +} diff --git a/flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory b/flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory similarity index 94% rename from flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory rename to flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory index 2b6bfa3cd579..037f55a5f0c5 100644 --- a/flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory +++ b/flink/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory @@ -13,4 +13,4 @@ # See the License for the specific language governing permissions and # limitations under the License. -org.apache.iceberg.flink.FlinkCatalogFactory +org.apache.iceberg.flink.FlinkCatalogFactory \ No newline at end of file diff --git a/versions.props b/versions.props index 0f28f66f8fe1..e7b001928beb 100644 --- a/versions.props +++ b/versions.props @@ -1,7 +1,7 @@ org.slf4j:* = 1.7.25 org.apache.avro:avro = 1.9.2 org.apache.calcite:* = 1.10.0 -org.apache.flink:* = 1.12.1 +org.apache.flink:* = 1.13.0 org.apache.hadoop:* = 2.7.3 org.apache.hive:hive-metastore = 2.3.8 org.apache.hive:hive-serde = 2.3.8 From bd36845f276ccfc99bad6d3ec5d64b80f5a8677e Mon Sep 17 00:00:00 2001 From: zhangjun30 Date: Mon, 24 May 2021 18:25:22 +0800 Subject: [PATCH 2/2] base-namespace --- .../main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java | 2 +- .../org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java index ff93bee78aaa..1d2ce74c0cc0 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java +++ b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactory.java @@ -118,7 +118,7 @@ public Set> optionalOptions() { options.add(FlinkCatalogFactoryOptions.URI); options.add(FlinkCatalogFactoryOptions.WAREHOUSE); options.add(FlinkCatalogFactoryOptions.CLIENTS); - options.add(FlinkCatalogFactoryOptions.BASE_NAMESPACES); + options.add(FlinkCatalogFactoryOptions.BASE_NAMESPACE); options.add(FlinkCatalogFactoryOptions.HIVE_CONF_DIF); options.add(FlinkCatalogFactoryOptions.CACHE_ENABLED); return options; diff --git a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java index 3ba913d22453..672e69a28e3d 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java +++ b/flink/src/main/java/org/apache/iceberg/flink/FlinkCatalogFactoryOptions.java @@ -64,7 +64,7 @@ private FlinkCatalogFactoryOptions() { " /hive-site.xml (or hive configure file from classpath) will be overwrote with " + "the warehouse value if setting both hive-conf-dir and warehouse when creating iceberg catalog."); - public static final ConfigOption BASE_NAMESPACES = + public static final ConfigOption BASE_NAMESPACE = ConfigOptions.key("base-namespace") .stringType() .noDefaultValue()