From ef19f26fd7bb17a363a9d137f8d23fa91347af27 Mon Sep 17 00:00:00 2001 From: Ankur Kotwal Date: Wed, 7 Oct 2026 11:23:42 -0700 Subject: [PATCH 1/4] Issue #8794 : Write Iceberg PATH tables at their path on the Spark engine An Iceberg PATH table was addressed as hop_iceberg.`` in a single Hadoop catalog whose warehouse is in java.io.tmpdir. Iceberg only treats an identifier as a location when Spark builds it from load(path), so the catalog took the whole URI as a table name and stored the table under its own warehouse instead of at the path. Each folder that holds PATH tables now gets its own Hadoop catalog, hop_iceberg_, with that folder as warehouse, and a table is addressed as hop_iceberg_.``. The Hadoop catalog keeps such a table at /, which is the table path. Reads, writes, merges, time travel and maintenance procedures all go through it. Generated-by: Claude Opus 5.5 --- .../pages/metadata-types/spark-catalog.adoc | 2 +- .../ROOT/pages/pipeline/spark/lakehouse.adoc | 10 ++- .../transforms/spark-lake-table-input.adoc | 2 +- .../spark-lake-table-maintenance.adoc | 2 +- .../transforms/spark-lake-table-merge.adoc | 2 +- .../transforms/spark-lake-table-output.adoc | 2 +- .../spark/table/SparkLakeTableSupport.java | 90 ++++++++++++++++--- .../table/SparkLakeTableIcebergPathTest.java | 8 ++ .../table/SparkLakeTableSupportTest.java | 44 +++++++-- 9 files changed, 134 insertions(+), 28 deletions(-) diff --git a/docs/hop-user-manual/modules/ROOT/pages/metadata-types/spark-catalog.adoc b/docs/hop-user-manual/modules/ROOT/pages/metadata-types/spark-catalog.adoc index 278f6ac2216..96ecd26dc28 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/metadata-types/spark-catalog.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/metadata-types/spark-catalog.adoc @@ -27,7 +27,7 @@ under the License. Lake Table transforms in **TABLE** mode can reference a Spark Catalog by name. Hop expands the entry into `spark.sql.catalog..*` configuration when building or preparing the session (`LakeSessionPlan` / `SparkCatalogApplier`). -PATH-mode Iceberg tables do **not** require this metadata type — the engine auto-registers a built-in Hadoop catalog named `hop_iceberg`. Use Spark Catalog for production Iceberg warehouses, REST catalogs, and other named catalogs. +PATH-mode Iceberg tables do **not** require this metadata type — the engine registers a Hadoop catalog for the folder that holds each table (`hop_iceberg_`). Use Spark Catalog for production Iceberg warehouses, REST catalogs, and other named catalogs. Connector packaging and session matrix: xref:pipeline/spark/lakehouse.adoc[Lakehouse tables on the native Spark engine]. diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc index 50588960a6a..6e217dda699 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc @@ -160,13 +160,15 @@ Hop’s `LakeSessionPlan` scans Lake Table transforms before the session is used |`spark.sql.extensions` includes `io.delta.sql.DeltaSparkSessionExtension` **and** `spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog` |Any Iceberg PATH use -|Iceberg extensions + built-in Hadoop catalog **`hop_iceberg`** (`type=hadoop`, warehouse under tmp) +|Iceberg extensions + a Hadoop catalog per table folder, **`hop_iceberg_`** (`type=hadoop`, warehouse = the folder that holds the table) |Iceberg TABLE |`spark.sql.catalog.=org.apache.iceberg.spark.SparkCatalog` (+ `type`, `warehouse` / `uri`, …) from xref:metadata-types/spark-catalog.adoc[Spark Catalog] metadata |=== -IMPORTANT: Delta PATH I/O on Spark 4.1 + Delta 4.3 **requires DeltaCatalog**, not only the extension. Iceberg bare `format("iceberg").load(path)` is **not** used — it defaults toward HiveCatalog and fails without HMS. PATH mode uses `hop_iceberg.\`file:///…\`` (or equivalent URI) instead. +IMPORTANT: Delta PATH I/O on Spark 4.1 + Delta 4.3 **requires DeltaCatalog**, not only the extension. Iceberg bare `format("iceberg").load(path)` is **not** used — it defaults toward HiveCatalog and fails without HMS. PATH mode uses `hop_iceberg_.\`\``, where the catalog's warehouse is the table's parent folder, so the table is read and written at its path. + +NOTE: Before Hop 2.20, Iceberg PATH tables written by the Spark engine were stored under `/hop-iceberg-path-catalog-warehouse/` instead of at the table path. To keep such a table, copy its folder from there (it is in a sub-folder named after the full table URI) to the table path before running the pipeline with Hop 2.20 or later. [[path-vs-table]] == PATH vs TABLE @@ -198,8 +200,8 @@ IMPORTANT: Delta PATH I/O on Spark 4.1 + Delta 4.3 **requires DeltaCatalog**, no | `format("delta").mode(…).save(path)` | **Iceberg** -| SQL over `hop_iceberg.uri` -| `writeTo(hop_iceberg.uri).using("iceberg")` +| SQL over `hop_iceberg_.table` +| `writeTo(hop_iceberg_.table).using("iceberg")` |=== diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc index 5e6901bf567..cffc18c66de 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc @@ -104,7 +104,7 @@ Typical pattern: template pipeline with Lake Table Input; injecting pipeline sup == Notes -* Prefer Lake Table Input over *Spark file input* for `delta` / `iceberg` so session extensions, `hop_iceberg` PATH catalog, and probes stay consistent. +* Prefer Lake Table Input over *Spark file input* for `delta` / `iceberg` so session extensions, Iceberg PATH catalogs, and probes stay consistent. * See the xref:pipeline/spark/lakehouse.adoc[lakehouse guide] for packaging, PATH vs TABLE, and troubleshooting. == Related pages diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc index ee7d966fab8..97029fc01f1 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc @@ -100,7 +100,7 @@ Aliases accepted for operation: `COMPACT` / `REWRITE_DATA_FILES` → OPTIMIZE; ` == Behaviour * SQL is executed as a Spark action during graph materialisation; an empty leaf Dataset is registered afterward. -* Iceberg CALL procedures use the catalog name from TABLE mode (`lake` in `lake.db.t`) or `hop_iceberg` for PATH mode. +* Iceberg CALL procedures use the catalog name from TABLE mode (`lake` in `lake.db.t`) or, for PATH mode, the `hop_iceberg_` catalog of the folder that holds the table. * Metrics are **log-only**; expect little or no GUI row traffic on this sink. == Metadata injection diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc index ab116c83d6b..a654871ab52 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc @@ -78,7 +78,7 @@ ON t.id = s.id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * ---- -* Target SQL ids: Delta PATH → `delta.\`path\``; Iceberg PATH → `hop_iceberg.\`uri\``; TABLE → multi-part identifier. +* Target SQL ids: Delta PATH → `delta.\`path\``; Iceberg PATH → `hop_iceberg_.\`table\`` (catalog warehouse = the table's parent folder); TABLE → multi-part identifier. * At least one action (matched and/or not-matched) is required. * Metrics are **log-only** (merge completion); GUI in/out row counts for the sink may stay at zero. diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc index b6fffd7355d..683fdd67890 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc @@ -72,7 +72,7 @@ This transform is **native Spark only**. Connectors ship with the native Spark p * The write runs as a Spark **action** when the engine materialises the graph. * An empty leaf Dataset is registered so a later engine-wide `count()` does not re-write the table. * **Delta PATH:** `format("delta").mode(…).save(path)` (requires Delta extension + `DeltaCatalog` on the session — hop-run applies these when Delta transforms are present). -* **Iceberg PATH:** `writeTo(hop_iceberg.\`uri\`).using("iceberg")` with the built-in `hop_iceberg` Hadoop catalog. +* **Iceberg PATH:** `writeTo(hop_iceberg_.\`table\`).using("iceberg")`, with a Hadoop catalog whose warehouse is the table's parent folder, so the table is written at its path. * **Iceberg TABLE:** `writeTo(catalog.ns.table).using("iceberg")`. * **Delta TABLE (advanced):** `format("delta").saveAsTable(id)`. diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java index 7cc0bf747b9..7ae035e9eea 100644 --- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java +++ b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java @@ -107,15 +107,63 @@ public static String toTableLocationUri(String path) { } /** - * Spark SQL multi-part identifier for an Iceberg path-based table under the hop path catalog: - * {@code hop_iceberg.`file:///path/to/table`}. + * An Iceberg table at a path, seen through the Hadoop catalog whose warehouse is the table's + * parent folder. A Hadoop catalog keeps a table of the default namespace in {@code + * /
}, so the table is read and written exactly at its path. + * + *

A quoted path in a single shared catalog, like {@code hop_iceberg.`file:///data/orders`}, + * doesn't work: Iceberg only treats an identifier as a location when Spark creates it from {@code + * DataFrameReader.load(path)}. Otherwise the Hadoop catalog takes the whole URI as a table name + * and stores the table under its own warehouse. */ - @SuppressWarnings("javabugs:S2259") // toTableLocationUri() never returns null for a resolved path + public record IcebergPathTable(String catalogName, String warehouse, String tableName) { + + /** Spark SQL identifier of the table, e.g. {@code hop_iceberg_1a2b3c4d5e6f.`orders`}. */ + public String sqlIdentifier() { + return catalogName + "." + procedureTableRef(); + } + + /** The table as named in an Iceberg procedure call of {@link #catalogName()}. */ + public String procedureTableRef() { + return "`" + tableName.replace("`", "``") + "`"; + } + } + + /** + * The path catalog and table name for an Iceberg table at {@code path}. Tables in the same folder + * share a catalog. + */ + public static IcebergPathTable icebergPathTable(String path) { + String uri = StringUtils.removeEnd(toTableLocationUri(path), "/"); + int slash = uri == null ? -1 : uri.lastIndexOf('/'); + String parent = slash < 0 ? "" : uri.substring(0, slash); + String name = slash < 0 ? "" : uri.substring(slash + 1); + if (name.isEmpty() || parent.isEmpty() || parent.endsWith(":") || parent.endsWith("/")) { + throw new IllegalArgumentException( + "Iceberg table path '" + path + "' needs a parent folder, e.g. file:///data/orders"); + } + return new IcebergPathTable( + SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + "_" + shortHash(parent), parent, name); + } + + private static String shortHash(String value) { + try { + byte[] digest = + java.security.MessageDigest.getInstance("SHA-256") + .digest(value.getBytes(java.nio.charset.StandardCharsets.UTF_8)); + StringBuilder hex = new StringBuilder(); + for (int i = 0; i < 6; i++) { + hex.append(String.format("%02x", digest[i])); + } + return hex.toString(); + } catch (java.security.NoSuchAlgorithmException e) { + throw new IllegalStateException(e); + } + } + + /** Spark SQL identifier for an Iceberg table at a path; see {@link IcebergPathTable}. */ public static String icebergPathSqlIdentifier(String path) { - String uri = toTableLocationUri(path); - // Escape any backticks in the URI (unlikely) by doubling them for Spark SQL quoting - String escaped = uri.replace("`", "``"); - return SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + ".`" + escaped + "`"; + return icebergPathTable(path).sqlIdentifier(); } public static String normalizeTimeTravelType(String type) throws HopException { @@ -256,10 +304,11 @@ public static MaintenanceTarget resolveMaintenanceTarget( String tableRefForCall = null; if (SparkLakeFormats.FORMAT_ICEBERG.equals(format)) { if (LakeTableInputMeta.MODE_PATH.equals(mode)) { - procedureCatalog = SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME; - tableRefForCall = - toTableLocationUri( + IcebergPathTable pathTable = + icebergPathTable( SparkPathDialect.toSparkUri(variables.resolve(meta.getTablePath()), pathSchemeMap)); + procedureCatalog = pathTable.catalogName(); + tableRefForCall = pathTable.procedureTableRef(); } else { String tableId = resolveTableIdentifier(meta.getTableIdentifier(), null, variables); String[] parts = tableId.split("\\."); @@ -338,7 +387,7 @@ public static String resolveLakeTargetSqlId( } if (spark != null) { - ensureIcebergPathCatalog(spark); + ensureIcebergPathCatalog(spark, path); } return icebergPathSqlIdentifier(path); } @@ -680,7 +729,7 @@ private static Dataset readIcebergPath( String timestamp, String transformName) throws HopException { - ensureIcebergPathCatalog(spark); + ensureIcebergPathCatalog(spark, path); String sqlId = icebergPathSqlIdentifier(path); String sql = buildIcebergTimeTravelSql(sqlId, timeTravelType, version, timestamp); try { @@ -741,7 +790,7 @@ private static void writeIcebergPath( String[] partitionColumns, String transformName) throws HopException { - ensureIcebergPathCatalog(spark); + ensureIcebergPathCatalog(spark, path); String sqlId = icebergPathSqlIdentifier(path); try { @@ -835,6 +884,21 @@ private static org.apache.spark.sql.DataFrameWriterV2 applyPartitioning( * Ensure the built-in Hadoop catalog for path identifiers is registered on this session. Safe to * call multiple times; no-ops when already present. */ + /** + * Registers the Hadoop catalog that serves the Iceberg table at {@code path} (see {@link + * IcebergPathTable}) on this session. Safe to call multiple times. + */ + public static void ensureIcebergPathCatalog(SparkSession spark, String path) { + IcebergPathTable table = icebergPathTable(path); + String key = "spark.sql.catalog." + table.catalogName(); + if (StringUtils.isNotEmpty(spark.conf().get(key, ""))) { + return; + } + spark.conf().set(key, SparkLakeFormats.ICEBERG_CATALOG); + spark.conf().set(key + ".type", "hadoop"); + spark.conf().set(key + ".warehouse", table.warehouse()); + } + public static void ensureIcebergPathCatalog(SparkSession spark) { String existing = spark.conf().get(SparkLakeFormats.SPARK_CONF_ICEBERG_PATH_CATALOG, ""); if (StringUtils.isNotEmpty(existing)) { diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java index d4b5e532475..658bba54be1 100644 --- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java +++ b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java @@ -18,8 +18,11 @@ package org.apache.hop.spark.table; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assumptions.assumeTrue; +import java.nio.file.Files; import java.nio.file.Path; import java.util.HashMap; import java.util.List; @@ -135,6 +138,11 @@ void icebergPathOutputThenInputRoundTrip() throws Exception { source); assertEquals(0, map.get("ice_out").count()); + // The table is written at its path, not under a catalog warehouse elsewhere. + assertTrue( + Files.exists(tablePath.resolve("metadata/version-hint.text")), + "Iceberg metadata is written under the table path"); + assertFalse(Files.exists(warehouse), "nothing is written to the hop_iceberg warehouse"); LakeTableInputMeta inMeta = new LakeTableInputMeta(); inMeta.setFormat(SparkLakeFormats.FORMAT_ICEBERG); diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java index 2470c825a30..d57e50fcd0d 100644 --- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java +++ b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java @@ -18,6 +18,7 @@ package org.apache.hop.spark.table; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -59,12 +60,43 @@ void collectFormatAddsNormalized() throws Exception { } @Test - void icebergPathSqlIdentifierQuotesUri() { - String id = SparkLakeTableSupport.icebergPathSqlIdentifier("/tmp/orders"); - assertTrue(id.startsWith(SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + ".`")); - assertTrue(id.endsWith("`")); - assertTrue(id.contains("file:")); - assertTrue(id.contains("orders")); + void icebergPathTableUsesTheParentFolderAsWarehouse() { + SparkLakeTableSupport.IcebergPathTable table = + SparkLakeTableSupport.icebergPathTable("s3a://bucket/lake/orders/"); + + assertEquals("s3a://bucket/lake", table.warehouse()); + assertEquals("orders", table.tableName()); + assertTrue(table.catalogName().startsWith(SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + "_")); + assertEquals(table.catalogName() + ".`orders`", table.sqlIdentifier()); + assertEquals("`orders`", table.procedureTableRef()); + } + + @Test + void icebergPathTablesShareACatalogPerFolder() { + String orders = + SparkLakeTableSupport.icebergPathTable("file:///data/lake/orders").catalogName(); + String items = SparkLakeTableSupport.icebergPathTable("file:///data/lake/items").catalogName(); + String other = + SparkLakeTableSupport.icebergPathTable("file:///data/other/orders").catalogName(); + + assertEquals(orders, items); + assertNotEquals(orders, other); + } + + @Test + void icebergPathSqlIdentifierQuotesTheTableName() { + String id = SparkLakeTableSupport.icebergPathSqlIdentifier("/tmp/my-orders"); + assertTrue(id.startsWith(SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + "_")); + assertTrue(id.endsWith(".`my-orders`")); + } + + @Test + void icebergPathTableNeedsAParentFolder() { + assertThrows( + IllegalArgumentException.class, + () -> SparkLakeTableSupport.icebergPathTable("s3a://bucket")); + assertThrows( + IllegalArgumentException.class, () -> SparkLakeTableSupport.icebergPathTable("file:///t")); } @Test From 3f519548dc9eec860a120c12d18bfe1dc1b4f971 Mon Sep 17 00:00:00 2001 From: Ankur Kotwal Date: Wed, 7 Oct 2026 23:34:44 -0700 Subject: [PATCH 2/4] Issue #8794 : Update the remaining hop_iceberg references for PATH tables Address review: put the Javadoc of both ensureIcebergPathCatalog methods back on the right method, describe PATH identifiers as hop_iceberg_.`table` in the class and resolver Javadoc, drop the hop_iceberg hint from the PATH read and write errors, and update the troubleshooting table in the lakehouse guide. Generated-by: Claude Opus 5.5 --- .../ROOT/pages/pipeline/spark/lakehouse.adoc | 2 +- .../hop/spark/table/SparkLakeFormats.java | 6 +++-- .../spark/table/SparkLakeTableSupport.java | 26 +++++++++---------- 3 files changed, 18 insertions(+), 16 deletions(-) diff --git a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc index 6e217dda699..e44db08e66b 100644 --- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc +++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/lakehouse.adoc @@ -292,7 +292,7 @@ Standalone *Spark lake table maintenance*: |`spark_catalog` must be `DeltaCatalog` when using Delta (hop-run sets this when Delta transforms are present) |Iceberg PATH fails with HiveCatalog / HMS -|Do not use bare `format("iceberg")` outside Hop; ensure Iceberg transforms run under hop-run so `hop_iceberg` is registered, or pass equivalent conf on submit +|Do not use bare `format("iceberg")` outside Hop; ensure Iceberg transforms run under hop-run so the per-folder `hop_iceberg_` catalogs are registered |TABLE unresolved |Spark Catalog metadata exported? Catalog name matches the first segment of `catalog.ns.table`? Warehouse / URI reachable from driver and executors? diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeFormats.java b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeFormats.java index d5a80b151e3..871658ce2ff 100644 --- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeFormats.java +++ b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeFormats.java @@ -48,8 +48,10 @@ public final class SparkLakeFormats { public static final String ICEBERG_CATALOG = LakeFormats.ICEBERG_CATALOG; /** - * Built-in Hadoop catalog name used for Iceberg PATH mode ({@code hop_iceberg.`file:///…`}). - * Distinct from {@code spark_catalog} so Delta can keep DeltaCatalog when both formats co-exist. + * Name prefix of the Hadoop catalogs used for Iceberg PATH mode: one catalog per table folder, + * for example {@code hop_iceberg_1a2b3c4d5e6f.`orders`} (see {@link + * SparkLakeTableSupport#icebergPathTable(String)}). Distinct from {@code spark_catalog} so Delta + * can keep DeltaCatalog when both formats co-exist. */ public static final String ICEBERG_PATH_CATALOG_NAME = "hop_iceberg"; diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java index 7ae035e9eea..2571560b4e2 100644 --- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java +++ b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java @@ -46,8 +46,9 @@ * *

    *
  • Delta PATH: {@code format("delta").load/save(path)} (requires DeltaCatalog on session) - *
  • Iceberg PATH: path identifier under built-in Hadoop catalog {@code hop_iceberg.`uri`} — - * bare {@code format("iceberg").load(path)} defaults to HiveCatalog and is not used + *
  • Iceberg PATH: {@code hop_iceberg_.`table`}, in a Hadoop catalog whose warehouse is + * the table's parent folder (see {@link IcebergPathTable}) — bare {@code + * format("iceberg").load(path)} defaults to HiveCatalog and is not used *
*/ public final class SparkLakeTableSupport { @@ -235,7 +236,7 @@ public static Map timeTravelOptionMap( * *
    *
  • Delta PATH → {@code delta.`path`} - *
  • Iceberg PATH → {@code hop_iceberg.`uri`} + *
  • Iceberg PATH → {@code hop_iceberg_.`table`} (see {@link IcebergPathTable}) *
  • TABLE → resolved multi-part table id *
*/ @@ -743,10 +744,8 @@ private static Dataset readIcebergPath( + sql + ") in transform '" + transformName - + "'. Ensure the Iceberg runtime is on the engine classpath and session has" - + " IcebergSparkSessionExtensions + Hadoop catalog '" - + SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME - + "' (see plugins/engines/spark/README.md).", + + "'. Ensure the Iceberg runtime is on the engine classpath and the session" + + " has IcebergSparkSessionExtensions (see plugins/engines/spark/README.md).", path), e); } @@ -856,8 +855,8 @@ private static void writeIcebergPath( + sqlId + ") in transform '" + transformName - + "'. Ensure Iceberg is on the classpath and hop_iceberg Hadoop catalog is" - + " configured (see plugins/engines/spark/README.md).", + + "'. Ensure Iceberg is on the classpath (see" + + " plugins/engines/spark/README.md).", path), e); } @@ -880,10 +879,6 @@ private static org.apache.spark.sql.DataFrameWriterV2 applyPartitioning( return writer.partitionedBy(first, rest); } - /** - * Ensure the built-in Hadoop catalog for path identifiers is registered on this session. Safe to - * call multiple times; no-ops when already present. - */ /** * Registers the Hadoop catalog that serves the Iceberg table at {@code path} (see {@link * IcebergPathTable}) on this session. Safe to call multiple times. @@ -899,6 +894,11 @@ public static void ensureIcebergPathCatalog(SparkSession spark, String path) { spark.conf().set(key + ".warehouse", table.warehouse()); } + /** + * Ensure the built-in {@code hop_iceberg} Hadoop catalog is registered on this session. It is + * used for two-part TABLE identifiers; PATH tables use {@link #ensureIcebergPathCatalog( + * SparkSession, String)}. Safe to call multiple times; no-ops when already present. + */ public static void ensureIcebergPathCatalog(SparkSession spark) { String existing = spark.conf().get(SparkLakeFormats.SPARK_CONF_ICEBERG_PATH_CATALOG, ""); if (StringUtils.isNotEmpty(existing)) { From 002f40e9111f22f155611ac3c6a2152662d5d0b2 Mon Sep 17 00:00:00 2001 From: Ankur Kotwal Date: Fri, 9 Oct 2026 10:37:27 -0700 Subject: [PATCH 3/4] Issue #8794 : One catalog per folder, whatever the spelling of the path Address review: the path of an Iceberg PATH table is canonicalized before it is split into warehouse and table name, and before the catalog name is hashed: the scheme is lowercased, file:/x, file://localhost/x and file:///x become file:///x, empty segments are collapsed and . and .. are resolved. Two spellings of one folder now share one catalog, so a write through one isn't hidden by the other's cached snapshot. Only a parent that is a scheme alone (file:, file://, s3a:/) is rejected now, so a table in a Windows drive root (file:///C:/orders) and paths with an empty segment (s3a://bucket/lake//orders) work. An invalid path is reported as a HopException that names the transform. Generated-by: Claude Opus 5.5 --- .../spark/table/SparkLakeTableSupport.java | 107 ++++++++++++++++-- .../table/SparkLakeTableSupportTest.java | 49 ++++++++ 2 files changed, 146 insertions(+), 10 deletions(-) diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java index 2571560b4e2..aa881622687 100644 --- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java +++ b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java @@ -23,6 +23,8 @@ import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import org.apache.commons.lang3.StringUtils; import org.apache.hop.core.Const; import org.apache.hop.core.exception.HopException; @@ -135,11 +137,11 @@ public String procedureTableRef() { * share a catalog. */ public static IcebergPathTable icebergPathTable(String path) { - String uri = StringUtils.removeEnd(toTableLocationUri(path), "/"); + String uri = canonicalLocation(toTableLocationUri(path)); int slash = uri == null ? -1 : uri.lastIndexOf('/'); String parent = slash < 0 ? "" : uri.substring(0, slash); String name = slash < 0 ? "" : uri.substring(slash + 1); - if (name.isEmpty() || parent.isEmpty() || parent.endsWith(":") || parent.endsWith("/")) { + if (name.isEmpty() || isSchemeOnly(parent)) { throw new IllegalArgumentException( "Iceberg table path '" + path + "' needs a parent folder, e.g. file:///data/orders"); } @@ -147,6 +149,83 @@ public static IcebergPathTable icebergPathTable(String path) { SparkLakeFormats.ICEBERG_PATH_CATALOG_NAME + "_" + shortHash(parent), parent, name); } + /** + * Like {@link #icebergPathTable(String)}, but reports an invalid path as a {@link HopException} + * that names the transform. + */ + public static IcebergPathTable icebergPathTable(String path, String transformName) + throws HopException { + try { + return icebergPathTable(path); + } catch (IllegalArgumentException e) { + throw new HopException("Spark Lake Table '" + transformName + "': " + e.getMessage(), e); + } + } + + private static final Pattern URI_SCHEME = Pattern.compile("^([A-Za-z][A-Za-z0-9+.-]*):(.*)$"); + + /** + * One spelling per location, so that one folder is always one catalog: the scheme is lowercased, + * {@code file:/x}, {@code file://localhost/x} and {@code file:///x} all become {@code file:///x}, + * empty path segments are collapsed, and {@code .} and {@code ..} are resolved. A trailing slash + * is removed. + */ + static String canonicalLocation(String uri) { + if (StringUtils.isEmpty(uri)) { + return uri; + } + Matcher matcher = URI_SCHEME.matcher(uri.trim()); + if (!matcher.matches() || matcher.group(1).length() == 1) { + // No scheme (or a Windows drive letter): leave it as it is. + return StringUtils.removeEnd(uri.trim(), "/"); + } + String scheme = matcher.group(1).toLowerCase(java.util.Locale.ROOT); + String rest = matcher.group(2); + String authority = ""; + String pathPart = rest; + if (rest.startsWith("//")) { + int end = rest.indexOf('/', 2); + authority = end < 0 ? rest.substring(2) : rest.substring(2, end); + pathPart = end < 0 ? "" : rest.substring(end); + } + if ("file".equals(scheme) && "localhost".equalsIgnoreCase(authority)) { + authority = ""; + } + java.util.Deque segments = new java.util.ArrayDeque<>(); + for (String segment : pathPart.split("/")) { + if (segment.isEmpty() || ".".equals(segment)) { + continue; + } + if ("..".equals(segment)) { + segments.pollLast(); + } else { + segments.addLast(segment); + } + } + String normalizedPath = segments.isEmpty() ? "" : "/" + String.join("/", segments); + return scheme + "://" + authority + normalizedPath; + } + + /** + * True for {@code file:}, {@code file://}, {@code s3a:/} and the like: a scheme and no folder. + */ + private static boolean isSchemeOnly(String parent) { + if (parent.isEmpty()) { + return true; + } + Matcher matcher = URI_SCHEME.matcher(parent); + if (!matcher.matches() || matcher.group(1).length() == 1) { + return false; + } + String rest = matcher.group(2); + if (!rest.startsWith("//")) { + return StringUtils.strip(rest, "/").isEmpty(); + } + // scheme://authority/path: a warehouse needs a bucket or host, or for file: a path. + String afterSlashes = rest.substring(2); + return StringUtils.strip(afterSlashes, "/").isEmpty(); + } + private static String shortHash(String value) { try { byte[] digest = @@ -307,7 +386,8 @@ public static MaintenanceTarget resolveMaintenanceTarget( if (LakeTableInputMeta.MODE_PATH.equals(mode)) { IcebergPathTable pathTable = icebergPathTable( - SparkPathDialect.toSparkUri(variables.resolve(meta.getTablePath()), pathSchemeMap)); + SparkPathDialect.toSparkUri(variables.resolve(meta.getTablePath()), pathSchemeMap), + transformName); procedureCatalog = pathTable.catalogName(); tableRefForCall = pathTable.procedureTableRef(); } else { @@ -387,10 +467,11 @@ public static String resolveLakeTargetSqlId( return "delta.`" + loc + "`"; } + IcebergPathTable pathTable = icebergPathTable(path, transformName); if (spark != null) { - ensureIcebergPathCatalog(spark, path); + ensureIcebergPathCatalog(spark, pathTable); } - return icebergPathSqlIdentifier(path); + return pathTable.sqlIdentifier(); } /** Maintenance target resolution result. */ @@ -730,8 +811,9 @@ private static Dataset readIcebergPath( String timestamp, String transformName) throws HopException { - ensureIcebergPathCatalog(spark, path); - String sqlId = icebergPathSqlIdentifier(path); + IcebergPathTable pathTable = icebergPathTable(path, transformName); + ensureIcebergPathCatalog(spark, pathTable); + String sqlId = pathTable.sqlIdentifier(); String sql = buildIcebergTimeTravelSql(sqlId, timeTravelType, version, timestamp); try { return spark.sql(sql); @@ -789,8 +871,9 @@ private static void writeIcebergPath( String[] partitionColumns, String transformName) throws HopException { - ensureIcebergPathCatalog(spark, path); - String sqlId = icebergPathSqlIdentifier(path); + IcebergPathTable pathTable = icebergPathTable(path, transformName); + ensureIcebergPathCatalog(spark, pathTable); + String sqlId = pathTable.sqlIdentifier(); try { switch (saveMode) { @@ -884,7 +967,11 @@ private static org.apache.spark.sql.DataFrameWriterV2 applyPartitioning( * IcebergPathTable}) on this session. Safe to call multiple times. */ public static void ensureIcebergPathCatalog(SparkSession spark, String path) { - IcebergPathTable table = icebergPathTable(path); + ensureIcebergPathCatalog(spark, icebergPathTable(path)); + } + + /** Registers the Hadoop catalog that serves {@code table} on this session. */ + public static void ensureIcebergPathCatalog(SparkSession spark, IcebergPathTable table) { String key = "spark.sql.catalog." + table.catalogName(); if (StringUtils.isNotEmpty(spark.conf().get(key, ""))) { return; diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java index d57e50fcd0d..309218f3372 100644 --- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java +++ b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java @@ -99,6 +99,55 @@ void icebergPathTableNeedsAParentFolder() { IllegalArgumentException.class, () -> SparkLakeTableSupport.icebergPathTable("file:///t")); } + @Test + void icebergPathTableAcceptsAWindowsDriveRoot() { + SparkLakeTableSupport.IcebergPathTable table = + SparkLakeTableSupport.icebergPathTable("file:///C:/orders"); + + assertEquals("file:///C:", table.warehouse()); + assertEquals("orders", table.tableName()); + } + + @Test + void icebergPathTableCollapsesEmptySegments() { + SparkLakeTableSupport.IcebergPathTable doubled = + SparkLakeTableSupport.icebergPathTable("s3a://bucket/lake//orders"); + SparkLakeTableSupport.IcebergPathTable single = + SparkLakeTableSupport.icebergPathTable("s3a://bucket/lake/orders"); + + assertEquals("s3a://bucket/lake", doubled.warehouse()); + assertEquals(single, doubled); + assertEquals( + single, SparkLakeTableSupport.icebergPathTable("s3a://bucket/lake/./tmp/../orders")); + } + + @Test + void oneFolderIsOneCatalogWhateverTheSpelling() { + SparkLakeTableSupport.IcebergPathTable canonical = + SparkLakeTableSupport.icebergPathTable("file:///data/lake/orders"); + + // File.toURI() and Path.toUri() spell the same folder differently. + for (String spelling : + new String[] { + "file:/data/lake/orders", + "file://localhost/data/lake/orders", + "FILE:///data/lake/orders/", + "file:///data//lake/orders" + }) { + assertEquals(canonical, SparkLakeTableSupport.icebergPathTable(spelling), spelling); + } + assertEquals("file:///data/lake", canonical.warehouse()); + } + + @Test + void invalidPathNamesTheTransform() { + HopException e = + assertThrows( + HopException.class, + () -> SparkLakeTableSupport.icebergPathTable("file:///t", "write orders")); + assertTrue(e.getMessage().contains("'write orders'"), e.getMessage()); + } + @Test void toTableLocationUriPreservesSchemes() { assertEquals("s3a://bucket/t", SparkLakeTableSupport.toTableLocationUri("s3a://bucket/t")); From 664246451ffef6a5cb8375060bf3db801272a45a Mon Sep 17 00:00:00 2001 From: Ankur Kotwal Date: Fri, 9 Oct 2026 12:15:10 -0700 Subject: [PATCH 4/4] Issue #8794 : Decode escapes and uppercase drive letters in PATH table locations Address second review. The canonical location is used as a Hadoop path, and Hadoop doesn't treat % as an escape, so a scheme-less path like /data/my lake/orders, which reaches us as file:///data/my%20lake/orders, was written to a folder literally named my%20lake. Percent-escapes in the path and authority are now decoded as UTF-8 (+ stays +, a segment with a broken escape is kept as is), and for file URIs a Windows drive letter is uppercased, so file:///c:/x and file:///C:/x are one catalog. A Spark round trip into a folder with spaces checks the table lands there; it failed before this change. Generated-by: Claude Opus 5.5 --- .../spark/table/SparkLakeTableSupport.java | 59 ++++++++++++++++++- .../table/SparkLakeTableIcebergPathTest.java | 17 +++++- .../table/SparkLakeTableSupportTest.java | 37 ++++++++++++ 3 files changed, 109 insertions(+), 4 deletions(-) diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java index aa881622687..ce3d3afa581 100644 --- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java +++ b/plugins/engines/spark/src/main/java/org/apache/hop/spark/table/SparkLakeTableSupport.java @@ -167,8 +167,11 @@ public static IcebergPathTable icebergPathTable(String path, String transformNam /** * One spelling per location, so that one folder is always one catalog: the scheme is lowercased, * {@code file:/x}, {@code file://localhost/x} and {@code file:///x} all become {@code file:///x}, - * empty path segments are collapsed, and {@code .} and {@code ..} are resolved. A trailing slash - * is removed. + * empty path segments are collapsed, {@code .} and {@code ..} are resolved, and a trailing slash + * is removed. Percent-escapes in the path and the authority are decoded, because the result is + * used as a Hadoop path, and Hadoop doesn't treat {@code %} as an escape: {@code + * file:///data/my%20lake} would otherwise be the folder {@code my%20lake}. A segment with a + * broken escape is kept as it is. For {@code file} URIs, a Windows drive letter is uppercased. */ static String canonicalLocation(String uri) { if (StringUtils.isEmpty(uri)) { @@ -188,11 +191,13 @@ static String canonicalLocation(String uri) { authority = end < 0 ? rest.substring(2) : rest.substring(2, end); pathPart = end < 0 ? "" : rest.substring(end); } + authority = percentDecode(authority); if ("file".equals(scheme) && "localhost".equalsIgnoreCase(authority)) { authority = ""; } java.util.Deque segments = new java.util.ArrayDeque<>(); - for (String segment : pathPart.split("/")) { + for (String raw : pathPart.split("/")) { + String segment = percentDecode(raw); if (segment.isEmpty() || ".".equals(segment)) { continue; } @@ -202,10 +207,58 @@ static String canonicalLocation(String uri) { segments.addLast(segment); } } + if ("file".equals(scheme) + && !segments.isEmpty() + && WINDOWS_DRIVE.matcher(segments.peekFirst()).matches()) { + segments.addFirst(segments.pollFirst().toUpperCase(java.util.Locale.ROOT)); + } String normalizedPath = segments.isEmpty() ? "" : "/" + String.join("/", segments); return scheme + "://" + authority + normalizedPath; } + private static final Pattern WINDOWS_DRIVE = Pattern.compile("[A-Za-z]:"); + + /** + * Decodes {@code %XX} escapes as UTF-8, leaving {@code +} alone. A value with a malformed escape + * is returned unchanged. + */ + static String percentDecode(String value) { + if (value.indexOf('%') < 0) { + return value; + } + java.io.ByteArrayOutputStream bytes = new java.io.ByteArrayOutputStream(); + for (int i = 0; i < value.length(); ) { + char c = value.charAt(i); + if (c == '%') { + if (i + 2 >= value.length()) { + return value; + } + int hi = Character.digit(value.charAt(i + 1), 16); + int lo = Character.digit(value.charAt(i + 2), 16); + if (hi < 0 || lo < 0) { + return value; + } + bytes.write(hi * 16 + lo); + i += 3; + } else { + int end = i + Character.charCount(value.codePointAt(i)); + byte[] encoded = value.substring(i, end).getBytes(java.nio.charset.StandardCharsets.UTF_8); + bytes.write(encoded, 0, encoded.length); + i = end; + } + } + java.nio.charset.CharsetDecoder decoder = + java.nio.charset.StandardCharsets.UTF_8 + .newDecoder() + .onMalformedInput(java.nio.charset.CodingErrorAction.REPORT) + .onUnmappableCharacter(java.nio.charset.CodingErrorAction.REPORT); + try { + return decoder.decode(java.nio.ByteBuffer.wrap(bytes.toByteArray())).toString(); + } catch (java.nio.charset.CharacterCodingException e) { + return value; + } + } + /** * True for {@code file:}, {@code file://}, {@code s3a:/} and the like: a scheme and no folder. */ diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java index 658bba54be1..bcd5a5d7773 100644 --- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java +++ b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableIcebergPathTest.java @@ -83,11 +83,26 @@ void stopSpark() { @Test void icebergPathOutputThenInputRoundTrip() throws Exception { + roundTrip(tempDir.resolve("orders_iceberg")); + } + + /** + * A folder with spaces: the path reaches Spark escaped ({@code my%20lake}), and Hadoop doesn't + * treat {@code %} as an escape, so without decoding the table would land in a folder literally + * named {@code my%20lake}. + */ + @Test + void icebergPathWithSpacesIsWrittenThere() throws Exception { + Path tablePath = tempDir.resolve("my lake").resolve("my orders"); + roundTrip(tablePath); + assertFalse(Files.exists(tempDir.resolve("my%20lake")), "no folder with an escaped name"); + } + + private void roundTrip(Path tablePath) throws Exception { assumeTrue( SparkLakeConnectorProbe.isIcebergPresent(SparkLakeConnectorProbe.class.getClassLoader()), "Iceberg connector not on classpath; connectors missing from test classpath"); - Path tablePath = tempDir.resolve("orders_iceberg"); Path warehouse = tempDir.resolve("iceberg_wh"); spark = SparkSession.builder() diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java index 309218f3372..28898995bed 100644 --- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java +++ b/plugins/engines/spark/src/test/java/org/apache/hop/spark/table/SparkLakeTableSupportTest.java @@ -139,6 +139,43 @@ void oneFolderIsOneCatalogWhateverTheSpelling() { assertEquals("file:///data/lake", canonical.warehouse()); } + @Test + void escapedAndUnescapedSpellingsAreOneCatalog() { + SparkLakeTableSupport.IcebergPathTable expected = + SparkLakeTableSupport.icebergPathTable("file:///data/my lake/my orders"); + + assertEquals("file:///data/my lake", expected.warehouse()); + assertEquals("my orders", expected.tableName()); + assertEquals( + expected, SparkLakeTableSupport.icebergPathTable("file:///data/my%20lake/my%20orders")); + // A scheme-less path goes through Path.toUri(), which escapes the spaces. + assertEquals( + expected.tableName(), + SparkLakeTableSupport.icebergPathTable("/data/my lake/my orders").tableName()); + assertTrue( + SparkLakeTableSupport.icebergPathTable("/data/my lake/my orders") + .warehouse() + .endsWith("/data/my lake")); + } + + @Test + void plusAndBrokenEscapesAreKept() { + assertEquals("a+b", SparkLakeTableSupport.percentDecode("a+b")); + assertEquals("100%", SparkLakeTableSupport.percentDecode("100%")); + assertEquals("x%zzy", SparkLakeTableSupport.percentDecode("x%zzy")); + assertEquals("caf\u00e9", SparkLakeTableSupport.percentDecode("caf%C3%A9")); + } + + @Test + void driveLetterCaseIsOneCatalog() { + assertEquals( + SparkLakeTableSupport.icebergPathTable("file:///C:/data/orders"), + SparkLakeTableSupport.icebergPathTable("file:///c:/data/orders")); + assertEquals( + "file:///C:/data", + SparkLakeTableSupport.icebergPathTable("file:///c:/data/orders").warehouse()); + } + @Test void invalidPathNamesTheTransform() { HopException e =