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 assemblies/debug/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -409,6 +409,12 @@
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-tech-lakehouse</artifactId>
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-tech-parquet</artifactId>
Expand Down
6 changes: 6 additions & 0 deletions assemblies/plugins/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -1767,6 +1767,12 @@
<version>${project.version}</version>
<type>zip</type>
</dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-tech-lakehouse</artifactId>
<version>${project.version}</version>
<type>zip</type>
</dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-tech-minio</artifactId>
Expand Down
8 changes: 8 additions & 0 deletions plugins/engines/spark/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,14 @@
<scope>provided</scope>
</dependency>

<!-- Lake table transform metadata and the catalog metadata type -->
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-tech-lakehouse</artifactId>
<version>${project.version}</version>
<scope>provided</scope>
</dependency>

<!-- Transform metadata for native shuffle-aware handlers (driver-side only) -->
<dependency>
<groupId>org.apache.hop</groupId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1118,7 +1118,7 @@ private SparkSession buildSparkSession() throws HopException {
}
}

// Lakehouse: Delta/Iceberg extensions, hop_iceberg PATH catalog, SparkCatalog metadata.
// Lakehouse: Delta/Iceberg extensions, hop_iceberg PATH catalog, LakeCatalog metadata.
// Applied after run-config sparkConfigs so lake defaults fill gaps; explicit run-config
// keys already set above win if the user overrode them (except catalog apply overwrites).
if (!lakePlan.isEmpty()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,13 @@
import org.apache.hop.core.logging.ILogChannel;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.lakehouse.transforms.LakeTableInputMeta;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.core.SparkNativeMetrics;
import org.apache.hop.spark.engines.ISparkPipelineEngineRunConfiguration;
import org.apache.hop.spark.table.SparkLakeTableSupport;
import org.apache.hop.spark.transforms.table.SparkLakeTableInputMeta;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -60,7 +60,7 @@ public void handleTransform(
Dataset<Row> input)
throws HopException {

SparkLakeTableInputMeta meta = new SparkLakeTableInputMeta();
LakeTableInputMeta meta = new LakeTableInputMeta();
loadTransformMetadata(meta, transformMeta, metadataProvider, pipelineMeta);

String pathSchemeMap = runConfiguration != null ? runConfiguration.getPathSchemeMap() : null;
Expand All @@ -71,8 +71,7 @@ public void handleTransform(
transformDatasetMap.put(transformMeta.getName(), dataset);

String target =
SparkLakeTableInputMeta.MODE_TABLE.equalsIgnoreCase(
String.valueOf(meta.getIdentifierMode()))
LakeTableInputMeta.MODE_TABLE.equalsIgnoreCase(String.valueOf(meta.getIdentifierMode()))
? "table=" + variables.resolve(meta.getTableIdentifier())
: "path=" + variables.resolve(meta.getTablePath());
log.logBasic(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.hop.core.logging.ILogChannel;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.lakehouse.transforms.LakeTableMaintenanceMeta;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
Expand All @@ -32,7 +33,6 @@
import org.apache.hop.spark.table.SparkLakeTableSupport;
import org.apache.hop.spark.table.SparkLakeTableSupport.MaintenanceTarget;
import org.apache.hop.spark.table.SparkMaintenanceSqlBuilder;
import org.apache.hop.spark.transforms.table.SparkLakeTableMaintenanceMeta;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -66,7 +66,7 @@ public void handleTransform(
throws HopException {

// Zero-input: ignore upstream if hop-connected (KD-20)
SparkLakeTableMaintenanceMeta meta = new SparkLakeTableMaintenanceMeta();
LakeTableMaintenanceMeta meta = new LakeTableMaintenanceMeta();
loadTransformMetadata(meta, transformMeta, metadataProvider, pipelineMeta);

String operation =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,14 +24,14 @@
import org.apache.hop.core.logging.ILogChannel;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.lakehouse.transforms.LakeTableMergeMeta;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.engines.ISparkPipelineEngineRunConfiguration;
import org.apache.hop.spark.table.SparkLakeActionSupport;
import org.apache.hop.spark.table.SparkLakeTableSupport;
import org.apache.hop.spark.table.SparkMergeSqlBuilder;
import org.apache.hop.spark.transforms.table.SparkLakeTableMergeMeta;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -71,7 +71,7 @@ public void handleTransform(
+ "' requires exactly one upstream Dataset (source rows for USING).");
}

SparkLakeTableMergeMeta meta = new SparkLakeTableMergeMeta();
LakeTableMergeMeta meta = new LakeTableMergeMeta();
loadTransformMetadata(meta, transformMeta, metadataProvider, pipelineMeta);

String pathSchemeMap = runConfiguration != null ? runConfiguration.getPathSchemeMap() : null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,15 +23,15 @@
import org.apache.hop.core.logging.ILogChannel;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.lakehouse.transforms.LakeTableInputMeta;
import org.apache.hop.lakehouse.transforms.LakeTableOutputMeta;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.core.SparkNativeMetrics;
import org.apache.hop.spark.engines.ISparkPipelineEngineRunConfiguration;
import org.apache.hop.spark.table.SparkLakeActionSupport;
import org.apache.hop.spark.table.SparkLakeTableSupport;
import org.apache.hop.spark.transforms.table.SparkLakeTableInputMeta;
import org.apache.hop.spark.transforms.table.SparkLakeTableOutputMeta;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -72,7 +72,7 @@ public void handleTransform(
+ " (disabled hops and disconnected transforms are not executed on native Spark).");
}

SparkLakeTableOutputMeta meta = new SparkLakeTableOutputMeta();
LakeTableOutputMeta meta = new LakeTableOutputMeta();
loadTransformMetadata(meta, transformMeta, metadataProvider, pipelineMeta);

Dataset<Row> toWrite = trackMetrics(input, transformMeta, SparkNativeMetrics.Role.OUTPUT);
Expand All @@ -82,8 +82,7 @@ public void handleTransform(
SparkLakeActionSupport.putEmptyLeaf(transformDatasetMap, transformMeta.getName(), spark);

String target =
SparkLakeTableInputMeta.MODE_TABLE.equalsIgnoreCase(
String.valueOf(meta.getIdentifierMode()))
LakeTableInputMeta.MODE_TABLE.equalsIgnoreCase(String.valueOf(meta.getIdentifierMode()))
? "table=" + variables.resolve(meta.getTableIdentifier())
: "path=" + variables.resolve(meta.getTablePath());
log.logBasic(
Expand Down
Loading
Loading