Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions common/src/java/org/apache/hadoop/hive/conf/HiveConf.java
Original file line number Diff line number Diff line change
Expand Up @@ -2154,6 +2154,12 @@ public static enum ConfVars {
"If this is set to true the URI for auth will have the default location masked with DEFAULT_TABLE_LOCATION"),
HIVE_ICEBERG_ALLOW_DATAFILES_IN_TABLE_LOCATION_ONLY("hive.iceberg.allow.datafiles.in.table.location.only", false,
"If this is set to true, then all the data files being read should be withing the table location"),
HIVE_ICEBERG_REST_SCAN_PLANNING_MODE(
"hive.iceberg.rest.scan-planning-mode", "client",
new StringSet("client", "server"),
"Hive-wide default for Iceberg REST catalog scan planning in Tez split generation, aligned with\n" +
"catalog property scan-planning-mode. In server mode, scan planning is done by Iceberg REST Server.\n" +
"In client mode (default), split generation is done by Hive."),
HIVE_USE_EXPLICIT_RCFILE_HEADER("hive.exec.rcfile.use.explicit.header", true,
"If this is set the header for RCFiles will simply be RCF. If this is not\n" +
"set the header will be that borrowed from sequence files, e.g. SEQ- followed\n" +
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -284,8 +284,8 @@ public void configureInputJobCredentials(TableDesc tableDesc, Map<String, String
try {
Table table =
IcebergVendedCredentialUtil.getTableWithVendedCredentials(tableDesc.getProperties(), conf);
String catalogName = tableDesc.getProperties().getProperty(InputFormatConfig.CATALOG_NAME);
IcebergVendedCredentialUtil.propagateToJob(table, catalogName, null, secrets, conf);
IcebergVendedCredentialUtil.propagateToJob(
table, tableDesc.getProperty(InputFormatConfig.CATALOG_NAME), null, secrets, conf);
} catch (NoSuchTableException ex) {
// Table may not exist yet for CTAS; credentials will not be available.
}
Expand Down Expand Up @@ -343,6 +343,10 @@ public void commitJob(JobContext originalContext) {
public void configureJobConf(TableDesc tableDesc, JobConf jobConf) {
setCommonJobConf(jobConf);
configureOutputTableJobConf(tableDesc, jobConf);
if (tableDesc != null && tableDesc.getProperties() != null) {
RestCatalogScanPlanningUtil.propagateCatalogPropertiesToJob(
conf, tableDesc.getProperty(InputFormatConfig.CATALOG_NAME), jobConf);
}
if (IcebergVendedCredentialUtil.requestsVendedCredentials(tableDesc.getProperties(), conf)) {
IcebergVendedCredentialUtil.refreshVendedCredentialsIfMissing(tableDesc, jobConf, conf);
}
Expand Down Expand Up @@ -382,7 +386,7 @@ private static void configureOutputTableJobConf(TableDesc tableDesc, JobConf job
tables = (tables == null) ? tableName : tables + TABLE_NAME_SEPARATOR + tableName;
jobConf.set(InputFormatConfig.OUTPUT_TABLES, tables);

String catalogName = tableDesc.getProperties().getProperty(InputFormatConfig.CATALOG_NAME);
String catalogName = tableDesc.getProperty(InputFormatConfig.CATALOG_NAME);
if (catalogName != null) {
jobConf.set(InputFormatConfig.TABLE_CATALOG_PREFIX + tableName, catalogName);
}
Expand Down Expand Up @@ -1772,6 +1776,9 @@ static void overlayTableProperties(Configuration configuration, TableDesc tableD
}
props.put(InputFormatConfig.PARTITION_SPEC, PartitionSpecParser.toJson(spec));

String catalogName = props.getProperty(InputFormatConfig.CATALOG_NAME);
RestCatalogScanPlanningUtil.propagateCatalogPropertiesToJob(configuration, catalogName, map);

// We need to remove this otherwise the job.xml will be invalid as column comments are separated with '\0' and
// the serialization utils fail to serialize this character
map.remove("columns.comments");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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)) {

Copy link
Copy Markdown
Member

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reload here means Catalogs.loadTable instead of the serialized SerializableTable, so we get a live RESTTable for POST /plan, only in getSplits when 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 hit loadTable once.

Table table = Catalogs.loadTable(conf);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes - Catalogs.loadTable creates a new REST catalog (and HTTP client) on each call, and we don’t close it on this path today. That’s the same pattern as other Catalogs.loadTable usages in the handler. Here it runs once per getSplits during split planning, not per input split. A reload only happens when shouldReloadForServerSideScanPlanning is true.

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.
Expand Down
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(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

please check visibility modifiers in public methods if the actually need thos

@difin difin Oct 6, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The 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 boolean isV2TableOrAbove(BinaryOperator<String> props)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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 isV2TableOrAbove overloads). BinaryOperator doesn’t seem to fit here: in isV2TableOrAbove it’s get-with-default; here we need put-if-absent.

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());
}
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import java.io.IOException;
import java.io.UncheckedIOException;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.ExecutorService;
import java.util.function.Consumer;
import org.apache.commons.lang3.StringUtils;
Expand Down Expand Up @@ -56,7 +55,6 @@
import org.apache.iceberg.expressions.Expression;
import org.apache.iceberg.expressions.Expressions;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.mr.Catalogs;
import org.apache.iceberg.mr.InputFormatConfig;
import org.apache.iceberg.mr.hive.HiveTableUtil;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
Expand Down Expand Up @@ -157,16 +155,18 @@ private static <T extends Scan<T, FileScanTask, CombinedScanTask>> T applyConfig
@Override
public List<InputSplit> getSplits(JobContext context) {
Configuration conf = context.getConfiguration();
Table table = Optional
.ofNullable(HiveTableUtil.deserializeTable(conf, conf.get(InputFormatConfig.TABLE_IDENTIFIER)))
.orElseGet(() -> {
Table tbl = Catalogs.loadTable(conf);
conf.set(InputFormatConfig.TABLE_IDENTIFIER, tbl.name());
// planning-local conf only (never shipped): for credential-vending catalogs the loaded
// table's FileIO carries secrets, which must not reach a serialized job configuration
conf.set(InputFormatConfig.SERIALIZED_TABLE_PREFIX + tbl.name(), SerializationUtil.serializeToBase64(tbl));
return tbl;
});
String tableIdentifier = conf.get(InputFormatConfig.TABLE_IDENTIFIER);
Table table = HiveTableUtil.resolveTableForScanPlanning(conf, tableIdentifier);
if (StringUtils.isBlank(tableIdentifier)) {
tableIdentifier = table.name();
conf.set(InputFormatConfig.TABLE_IDENTIFIER, tableIdentifier);
}
String serializedTableKey = InputFormatConfig.SERIALIZED_TABLE_PREFIX + tableIdentifier;
if (conf.get(serializedTableKey) == null) {
// planning-local conf only (never shipped): for credential-vending catalogs the loaded
// table's FileIO carries secrets, which must not reach a serialized job configuration
conf.set(serializedTableKey, SerializationUtil.serializeToBase64(table));
}
final ExecutorService workerPool =
ThreadPools.newFixedThreadPool("iceberg-plan-worker-pool",
conf.getInt(SystemConfigs.WORKER_THREAD_POOL_SIZE.propertyKey(), ThreadPools.WORKER_THREAD_POOL_SIZE));
Expand Down
Loading
Loading