Repository navigation
HIVE-30056: Iceberg: Add client side support for REST catalog server-side scan planning #6789
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -53,6 +53,7 @@ | |
| import org.apache.hadoop.mapreduce.JobID; | ||
| import org.apache.iceberg.AppendFiles; | ||
| import org.apache.iceberg.BaseTable; | ||
| import org.apache.iceberg.CatalogUtil; | ||
| import org.apache.iceberg.DataFile; | ||
| import org.apache.iceberg.DeleteFiles; | ||
| import org.apache.iceberg.MetricsConfig; | ||
|
|
@@ -69,6 +70,7 @@ | |
| import org.apache.iceberg.hadoop.HadoopConfigurable; | ||
| import org.apache.iceberg.hadoop.HadoopFileIO; | ||
| import org.apache.iceberg.hadoop.Util; | ||
| import org.apache.iceberg.hive.IcebergCatalogProperties; | ||
| import org.apache.iceberg.io.FileIO; | ||
| import org.apache.iceberg.io.OutputFile; | ||
| import org.apache.iceberg.mapping.NameMapping; | ||
|
|
@@ -259,6 +261,71 @@ public static Table deserializeTable(Configuration config, String name) { | |
| return table; | ||
| } | ||
|
|
||
| /** | ||
| * Resolves the Iceberg {@link Table} for split generation ({@code IcebergInputFormat#getSplits}). | ||
| * | ||
| * <p>Serialized tables ({@link SerializableTable}) only carry a metadata snapshot and produce | ||
| * {@link org.apache.iceberg.DataTableScan} (client-side manifest planning). When REST catalog | ||
| * server-side scan planning is enabled, reload the live table from the catalog so | ||
| * {@code table.newScan()} returns {@link org.apache.iceberg.rest.RESTTableScan} and issues | ||
| * {@code POST /plan} on the REST server. | ||
| * | ||
| * <p>Intra-transaction read-after-write ({@link InputFormatConfig#TABLE_METADATA_LOCATION}) still | ||
| * uses the deserialized snapshot so uncommitted metadata is visible. | ||
| */ | ||
| public static Table resolveTableForScanPlanning(Configuration conf, String tableIdentifier) { | ||
| if (shouldReloadForServerSideScanPlanning(conf)) { | ||
| Table table = Catalogs.loadTable(conf); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I believe this will create a new RestCatalog Object/HttpClient on every call, curious if this can cause any connection leak?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes - |
||
| checkAndSetIoConfig(conf, table); | ||
| IcebergVendedCredentialUtil.applyFromJobConf(table, resolveCatalogName(conf), conf); | ||
| return table; | ||
| } | ||
|
|
||
| Table table = deserializeTable(conf, tableIdentifier); | ||
| if (table == null) { | ||
| table = Catalogs.loadTable(conf); | ||
| checkAndSetIoConfig(conf, table); | ||
| } | ||
| return table; | ||
| } | ||
|
|
||
| private static boolean shouldReloadForServerSideScanPlanning(Configuration conf) { | ||
| if (StringUtils.isNotBlank(conf.get(InputFormatConfig.TABLE_METADATA_LOCATION))) { | ||
| return false; | ||
| } | ||
| String catalogName = resolveCatalogName(conf); | ||
| if (StringUtils.isBlank(catalogName)) { | ||
| return false; | ||
| } | ||
| if (!CatalogUtil.ICEBERG_CATALOG_TYPE_REST.equals( | ||
| IcebergCatalogProperties.getCatalogType(conf, catalogName))) { | ||
| return false; | ||
| } | ||
| return RestCatalogScanPlanningUtil.isServerSidePlanningEnabled(catalogName, conf); | ||
| } | ||
|
|
||
| /** | ||
| * Resolves the Iceberg catalog name from {@link InputFormatConfig#CATALOG_NAME} on the | ||
| * configuration, or the session default ({@link IcebergCatalogProperties#getCatalogName(Configuration)}). | ||
| */ | ||
| public static String resolveCatalogName(Configuration conf) { | ||
| return resolveCatalogName(conf, conf != null ? conf.get(InputFormatConfig.CATALOG_NAME) : null); | ||
| } | ||
|
|
||
| /** | ||
| * Resolves the Iceberg catalog name from a per-table {@link InputFormatConfig#CATALOG_NAME} value | ||
| * (e.g. {@code TableDesc} properties) or the session default when that value is absent. | ||
| */ | ||
| public static String resolveCatalogName(Configuration conf, String catalogNameFromTable) { | ||
| if (StringUtils.isNotBlank(catalogNameFromTable)) { | ||
| return catalogNameFromTable; | ||
| } | ||
| if (conf == null) { | ||
| return null; | ||
| } | ||
| return IcebergCatalogProperties.getCatalogName(conf); | ||
| } | ||
|
|
||
| /** | ||
| * If enabled, it populates the FileIO's hadoop configuration with the input config object. | ||
| * This might be necessary when the table object was serialized without the FileIO config. | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,206 @@ | ||
| /* | ||
| * 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.mr.hive; | ||
|
|
||
| import java.util.Map; | ||
| import java.util.function.BiConsumer; | ||
| import org.apache.commons.lang3.StringUtils; | ||
| import org.apache.hadoop.conf.Configuration; | ||
| import org.apache.hadoop.hive.conf.HiveConf; | ||
| import org.apache.hadoop.hive.metastore.conf.MetastoreConf; | ||
| import org.apache.iceberg.CatalogUtil; | ||
| import org.apache.iceberg.hive.IcebergCatalogProperties; | ||
| import org.apache.iceberg.rest.RESTCatalogProperties; | ||
|
|
||
| /** | ||
| * Utilities for Iceberg REST catalog server-side scan planning configuration. | ||
| * | ||
| * <p>When a REST catalog server advertises the scan-planning endpoints and | ||
| * {@link RESTCatalogProperties#SCAN_PLANNING_MODE} is set to | ||
| * {@link RESTCatalogProperties.ScanPlanningMode#SERVER}, Iceberg's {@code RESTSessionCatalog} returns a | ||
| * {@code RESTTable} that delegates {@code planTasks()} to the server. Hive's | ||
| * {@code IcebergInputFormat} calls {@code scan.planTasks()} via | ||
| * {@link HiveTableUtil#resolveTableForScanPlanning}, which reloads the live REST catalog table | ||
| * (instead of a serialized metadata snapshot) when server mode is enabled and | ||
| * {@link HiveConf.ConfVars#HIVE_ICEBERG_REST_SCAN_PLANNING_MODE} is {@code server}. | ||
| * Operators can use this helper or set catalog {@code scan-planning-mode} directly in {@code hive-site.xml}. | ||
| * | ||
| * <p>Tests: {@code TestRestCatalogScanPlanningUtil} in {@code iceberg-handler}; | ||
| * {@code TestHiveIcebergServerSideScanPlanning} in {@code iceberg-handler}; embedded REST server | ||
| * tests {@code TestRestCatalogScanPlanningServerIT} and {@code TestHiveIcebergServerSideScanPlanningServerIT} | ||
| * in {@code itests/hive-iceberg-rest-server}. | ||
| * | ||
| * <p>{@link #setCatalogMode(Configuration, String, String)}, | ||
| * {@link #isCatalogServerMode(Configuration, String)}, and | ||
| * {@link #setHiveMode(Configuration, String)} are public so tests in other Maven modules | ||
| * (for example {@code itests/hive-iceberg-rest-server}) can configure scan planning; those modules | ||
| * compile against this artifact as a JAR and cannot call package-private members. They are not | ||
| * intended as a general operator or application API. | ||
| * | ||
| * @see <a href="https://iceberg.apache.org/docs/latest/catalog-properties/">REST catalog properties</a> | ||
| */ | ||
| public final class RestCatalogScanPlanningUtil { | ||
|
|
||
| private RestCatalogScanPlanningUtil() { | ||
| } | ||
|
|
||
| private static void setCatalogMode( | ||
| Configuration conf, String catalogName, RESTCatalogProperties.ScanPlanningMode mode) { | ||
| conf.set( | ||
| IcebergCatalogProperties.catalogPropertyConfigKey( | ||
| catalogName, RESTCatalogProperties.SCAN_PLANNING_MODE), | ||
| mode.modeName()); | ||
| } | ||
|
|
||
| public static void setCatalogMode(Configuration conf, String catalogName, String mode) { | ||
| setCatalogMode(conf, catalogName, RESTCatalogProperties.ScanPlanningMode.fromString(mode)); | ||
| } | ||
|
|
||
| static RESTCatalogProperties.ScanPlanningMode getCatalogMode( | ||
| Configuration conf, String catalogName) { | ||
| String mode = conf.get( | ||
| IcebergCatalogProperties.catalogPropertyConfigKey( | ||
| catalogName, RESTCatalogProperties.SCAN_PLANNING_MODE), | ||
| RESTCatalogProperties.SCAN_PLANNING_MODE_DEFAULT.modeName()); | ||
| return RESTCatalogProperties.ScanPlanningMode.fromString(mode); | ||
| } | ||
|
|
||
| public static boolean isCatalogServerMode(Configuration conf, String catalogName) { | ||
| return getCatalogMode(conf, catalogName) == RESTCatalogProperties.ScanPlanningMode.SERVER; | ||
| } | ||
|
|
||
| /** | ||
| * Returns true when Hive server-side REST scan planning is enabled in configuration. | ||
| */ | ||
| static boolean isHiveServerMode(Configuration conf) { | ||
| if (conf == null) { | ||
| return false; | ||
| } | ||
| return RESTCatalogProperties.ScanPlanningMode.SERVER == | ||
| RESTCatalogProperties.ScanPlanningMode.fromString(getHiveMode(conf)); | ||
| } | ||
|
|
||
| static String getHiveMode(Configuration conf) { | ||
| if (conf == null) { | ||
| return RESTCatalogProperties.SCAN_PLANNING_MODE_DEFAULT.modeName(); | ||
| } | ||
| return HiveConf.getVar(conf, HiveConf.ConfVars.HIVE_ICEBERG_REST_SCAN_PLANNING_MODE); | ||
| } | ||
|
|
||
| public static void setHiveMode(Configuration conf, String mode) { | ||
| HiveConf.setVar( | ||
| conf, | ||
| HiveConf.ConfVars.HIVE_ICEBERG_REST_SCAN_PLANNING_MODE, | ||
| RESTCatalogProperties.ScanPlanningMode.fromString(mode).modeName()); | ||
| } | ||
|
|
||
| /** | ||
| * Returns true when the catalog is configured for server-side scan planning and the Hive feature flag is on. | ||
| */ | ||
| static boolean isServerSidePlanningEnabled(String catalogName, Configuration conf) { | ||
| if (conf == null || StringUtils.isEmpty(catalogName)) { | ||
| return false; | ||
| } | ||
| return isHiveServerMode(conf) && isCatalogServerMode(conf, catalogName); | ||
| } | ||
|
|
||
| /** | ||
| * Returns true when catalog properties should be copied into the Tez job configuration so | ||
| * executors can reload a live REST catalog table for server-side scan planning. | ||
| */ | ||
| static boolean shouldPropagateCatalogPropertiesToJob(String catalogName, Configuration conf) { | ||
| String resolvedCatalogName = HiveTableUtil.resolveCatalogName(conf, catalogName); | ||
| if (StringUtils.isEmpty(resolvedCatalogName) || conf == null) { | ||
| return false; | ||
| } | ||
| if (!CatalogUtil.ICEBERG_CATALOG_TYPE_REST.equals( | ||
| IcebergCatalogProperties.getCatalogType(conf, resolvedCatalogName))) { | ||
| return false; | ||
| } | ||
| return isServerSidePlanningEnabled(resolvedCatalogName, conf); | ||
| } | ||
|
|
||
| /** | ||
| * Copies {@code iceberg.catalog.<catalog>.*} entries from the HS2 session configuration into Tez | ||
| * job properties so executors can reload a live REST catalog table for server-side scan planning. | ||
| * | ||
| * <p>Session-level {@code SET} commands and {@code hive-site.xml} catalog settings are not | ||
| * automatically present in the job configuration; without this step split generation falls back to | ||
| * the serialized metadata snapshot ({@code DataTableScan}). | ||
| */ | ||
| public static void propagateCatalogPropertiesToJob( | ||
| Configuration sessionConf, String catalogName, Map<String, String> jobProperties) { | ||
| if (sessionConf == null || jobProperties == null) { | ||
| return; | ||
| } | ||
| propagateCatalogPropertiesToJob(sessionConf, catalogName, jobProperties::putIfAbsent); | ||
| } | ||
|
|
||
| /** | ||
| * Copies REST catalog configuration from the HS2 session into a runtime job {@link Configuration}. | ||
| * Called from {@code configureJobConf} so split generation sees catalog URI/type/scan-planning-mode | ||
| * even when job properties were not copied yet. | ||
| */ | ||
| public static void propagateCatalogPropertiesToJob( | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. please check visibility modifiers in public methods if the actually need thos
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. there were a few methods that should stay public to be visible for testing from another package (VisibleForTesting annotation is not suitable here because the style check requires visibility to be package-private on methods annotated with VisibleForTesting), and few other public methods whose visibility could be reduced - done.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. looks like duplicate. maybe use BinaryOperator to hadle both Map and jobConf, see
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Addressed the duplication with a single private path and thin Map / JobConf overloads (BiConsumer, same idea as the two |
||
| Configuration sessionConf, String catalogName, Configuration jobConf) { | ||
| if (sessionConf == null || jobConf == null) { | ||
| return; | ||
| } | ||
| propagateCatalogPropertiesToJob( | ||
| sessionConf, | ||
| catalogName, | ||
| (key, value) -> { | ||
| if (jobConf.get(key) == null) { | ||
| jobConf.set(key, value); | ||
| } | ||
| }); | ||
| } | ||
|
|
||
| private static void propagateCatalogPropertiesToJob( | ||
| Configuration sessionConf, String catalogName, BiConsumer<String, String> jobPropertySink) { | ||
| if (!shouldPropagateCatalogPropertiesToJob(catalogName, sessionConf)) { | ||
| return; | ||
| } | ||
|
|
||
| String resolvedCatalogName = HiveTableUtil.resolveCatalogName(sessionConf, catalogName); | ||
| if (StringUtils.isEmpty(resolvedCatalogName)) { | ||
| return; | ||
| } | ||
|
|
||
| jobPropertySink.accept( | ||
| HiveConf.ConfVars.HIVE_ICEBERG_REST_SCAN_PLANNING_MODE.varname, | ||
| getHiveMode(sessionConf)); | ||
|
|
||
| String sessionDefaultCatalog = | ||
| MetastoreConf.getVar(sessionConf, MetastoreConf.ConfVars.CATALOG_DEFAULT); | ||
| if (StringUtils.isNotBlank(sessionDefaultCatalog)) { | ||
| jobPropertySink.accept(MetastoreConf.ConfVars.CATALOG_DEFAULT.getVarname(), sessionDefaultCatalog); | ||
| } | ||
|
|
||
| String catalogPrefix = | ||
| IcebergCatalogProperties.CATALOG_CONFIG_PREFIX + resolvedCatalogName + "."; | ||
| sessionConf.forEach( | ||
| entry -> { | ||
| if (entry.getKey().startsWith(catalogPrefix)) { | ||
| jobPropertySink.accept(entry.getKey(), entry.getValue()); | ||
| } | ||
| }); | ||
| } | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
what is this reload, isn't it expensive? what if we have multiple requests in same query/session?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Reload here means
Catalogs.loadTableinstead of the serializedSerializableTable, so we get a liveRESTTablefor POST /plan, only ingetSplitswhen REST + server scan mode are on (skipped for intra-txn TABLE_METADATA_LOCATION). It’s extra REST metadata cost vs deserialize, but it runs at split-planning time, not per task/row; multiple scans of the same table in one query can each hitloadTableonce.