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
Original file line number Diff line number Diff line change
Expand Up @@ -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.<catalogName>.*` 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_<hash>`). 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].

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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_<hash>`** (`type=hadoop`, warehouse = the folder that holds the table)

|Iceberg TABLE
|`spark.sql.catalog.<name>=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_<hash>.\`<table>\``, 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 `<java.io.tmpdir>/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
Expand Down Expand Up @@ -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_<hash>.table`
| `writeTo(hop_iceberg_<hash>.table).using("iceberg")`

|===

Expand Down Expand Up @@ -290,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_<hash>` 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?
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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_<hash>` catalog of the folder that holds the table.
* Metrics are **log-only**; expect little or no GUI row traffic on this sink.

== Metadata injection
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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_<hash>.\`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.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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_<hash>.\`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)`.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,9 @@
*
* <ul>
* <li>Delta PATH: {@code format("delta").load/save(path)} (requires DeltaCatalog on session)
* <li>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
* <li>Iceberg PATH: {@code hop_iceberg_<hash>.`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
* </ul>
*/
public final class SparkLakeTableSupport {
Expand Down Expand Up @@ -106,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
* <warehouse>/<table>}, so the table is read and written exactly at its path.
*
* <p>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 {
Expand Down Expand Up @@ -186,7 +235,7 @@ public static Map<String, String> timeTravelOptionMap(
*
* <ul>
* <li>Delta PATH → {@code delta.`path`}
* <li>Iceberg PATH → {@code hop_iceberg.`uri`}
* <li>Iceberg PATH → {@code hop_iceberg_<hash>.`table`} (see {@link IcebergPathTable})
* <li>TABLE → resolved multi-part table id
* </ul>
*/
Expand Down Expand Up @@ -255,10 +304,11 @@ public static MaintenanceTarget resolveMaintenanceTarget(
String tableRefForCall = null;
if (SparkLakeFormats.FORMAT_ICEBERG.equals(format)) {
if (SparkLakeTableInputMeta.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("\\.");
Expand Down Expand Up @@ -337,7 +387,7 @@ public static String resolveLakeTargetSqlId(
}

if (spark != null) {
ensureIcebergPathCatalog(spark);
ensureIcebergPathCatalog(spark, path);
}
return icebergPathSqlIdentifier(path);
}
Expand Down Expand Up @@ -666,7 +716,7 @@ private static Dataset<Row> readIcebergPath(
String timestamp,
String transformName)
throws HopException {
ensureIcebergPathCatalog(spark);
ensureIcebergPathCatalog(spark, path);
String sqlId = icebergPathSqlIdentifier(path);
String sql = buildIcebergTimeTravelSql(sqlId, timeTravelType, version, timestamp);
try {
Expand All @@ -680,10 +730,8 @@ private static Dataset<Row> 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);
}
Expand Down Expand Up @@ -727,7 +775,7 @@ private static void writeIcebergPath(
String[] partitionColumns,
String transformName)
throws HopException {
ensureIcebergPathCatalog(spark);
ensureIcebergPathCatalog(spark, path);
String sqlId = icebergPathSqlIdentifier(path);

try {
Expand Down Expand Up @@ -793,8 +841,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);
}
Expand All @@ -818,8 +866,24 @@ private static org.apache.spark.sql.DataFrameWriterV2<Row> 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());
}

/**
* 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, "");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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");

SparkLakeTableInputMeta inMeta = new SparkLakeTableInputMeta();
inMeta.setFormat(SparkLakeFormats.FORMAT_ICEBERG);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -60,12 +61,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
Expand Down
Loading