manager, LakeCatalog metadata) {
super(hopGui, manager, metadata);
}
@@ -138,13 +138,13 @@ public void afterButtonPressed(Object sourceObject) {
@Override
public void setWidgetsContent() {
- SparkCatalog meta = this.getMetadata();
+ LakeCatalog meta = this.getMetadata();
wName.setText(Const.NVL(meta.getName(), ""));
guiCompositeWidgets.setWidgetsContents(metadata, wWidgetsComposite, GUI_WIDGETS_PARENT_ID);
}
@Override
- public void getWidgetsContent(SparkCatalog meta) {
+ public void getWidgetsContent(LakeCatalog meta) {
meta.setName(wName.getText());
guiCompositeWidgets.getWidgetsContents(metadata, GUI_WIDGETS_PARENT_ID);
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplate.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplate.java
similarity index 85%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplate.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplate.java
index adb989395d4..92711ed4afd 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplate.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplate.java
@@ -15,28 +15,28 @@
* limitations under the License.
*/
-package org.apache.hop.spark.metadata.template;
+package org.apache.hop.lakehouse.metadata.template;
import java.util.Arrays;
import java.util.Objects;
import org.apache.commons.lang3.StringUtils;
-import org.apache.hop.spark.metadata.SparkCatalog;
-import org.apache.hop.spark.table.SparkLakeFormats;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.metadata.LakeCatalog;
/**
- * Named presets for {@link SparkCatalog} fields. Pure data — no SWT. Apply via {@link
- * #applyTo(SparkCatalog)}.
+ * Named presets for {@link LakeCatalog} fields. Pure data — no SWT. Apply via {@link
+ * #applyTo(LakeCatalog)}.
*
* Advanced presets include a commented {@code # docs: …} line in conf extra (ignored by {@code
* SparkCatalogApplier}) so operators can open vendor documentation without leaving Hop.
*/
-public enum SparkCatalogTemplate {
+public enum LakeCatalogTemplate {
ICEBERG_HADOOP_LOCAL(
"Iceberg Hadoop (local)",
"Named Iceberg Hadoop catalog on a local warehouse path. Use TABLE mode as lake.db.table.",
"lake",
- SparkCatalog.TYPE_HADOOP,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_HADOOP,
+ LakeFormats.ICEBERG_CATALOG,
"file:///tmp/hop-warehouse",
"",
confWithDocs(Docs.ICEBERG_SPARK, "")),
@@ -44,8 +44,8 @@ public enum SparkCatalogTemplate {
"Iceberg Hadoop (object store)",
"Iceberg Hadoop warehouse on object storage (replace bucket). Credentials via env/Hadoop conf.",
"lake",
- SparkCatalog.TYPE_HADOOP,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_HADOOP,
+ LakeFormats.ICEBERG_CATALOG,
"s3a://bucket/warehouse",
"",
confWithDocs(Docs.ICEBERG_AWS, "io-impl=org.apache.iceberg.aws.s3.S3FileIO")),
@@ -53,8 +53,8 @@ public enum SparkCatalogTemplate {
"Iceberg REST",
"Iceberg REST catalog. Replace the URI with your catalog service endpoint.",
"lake",
- SparkCatalog.TYPE_REST,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_REST,
+ LakeFormats.ICEBERG_CATALOG,
"",
"https://catalog.example.com/v1",
confWithDocs(Docs.ICEBERG_REST, "")),
@@ -62,8 +62,8 @@ public enum SparkCatalogTemplate {
"Iceberg REST (authenticated)",
"Iceberg REST with token auth. Put the secret in Credential (maps to catalog .token).",
"lake",
- SparkCatalog.TYPE_REST,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_REST,
+ LakeFormats.ICEBERG_CATALOG,
"",
"https://catalog.example.com/v1",
confWithDocs(
@@ -73,8 +73,8 @@ public enum SparkCatalogTemplate {
"Hive Metastore (advanced)",
"Iceberg via Hive Metastore. Requires Hive/Iceberg deps on the cluster (not packaged by default).",
"hive",
- SparkCatalog.TYPE_HIVE,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_HIVE,
+ LakeFormats.ICEBERG_CATALOG,
"",
"",
confWithDocs(Docs.ICEBERG_HIVE, "uri=thrift://hive-metastore:9083")),
@@ -82,8 +82,8 @@ public enum SparkCatalogTemplate {
"AWS Glue (advanced)",
"Iceberg via AWS Glue. Requires AWS/Iceberg deps and IAM on the cluster (not packaged by default).",
"glue",
- SparkCatalog.TYPE_GLUE,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_GLUE,
+ LakeFormats.ICEBERG_CATALOG,
"s3a://bucket/warehouse",
"",
confWithDocs(
@@ -93,8 +93,8 @@ public enum SparkCatalogTemplate {
"Nessie (advanced)",
"Iceberg Nessie catalog. Requires Nessie/Iceberg deps; replace URI, ref, and warehouse.",
"nessie",
- SparkCatalog.TYPE_CUSTOM,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_CUSTOM,
+ LakeFormats.ICEBERG_CATALOG,
"file:///tmp/nessie-warehouse",
"http://localhost:19120/api/v1",
confWithDocs(
@@ -107,8 +107,8 @@ public enum SparkCatalogTemplate {
"Databricks Unity Catalog (advanced)",
"Skeleton for Unity Catalog on Databricks. Prefer workspace-managed conf; adjust URI/auth for your env.",
"unity",
- SparkCatalog.TYPE_CUSTOM,
- SparkLakeFormats.ICEBERG_CATALOG,
+ LakeCatalog.TYPE_CUSTOM,
+ LakeFormats.ICEBERG_CATALOG,
"",
"",
confWithDocs(
@@ -121,8 +121,8 @@ public enum SparkCatalogTemplate {
"Delta named catalog (advanced)",
"Named DeltaCatalog (not spark_catalog). Prefer spark_catalog for Delta when coexisting with Iceberg.",
"delta_cat",
- SparkCatalog.TYPE_CUSTOM,
- SparkLakeFormats.DELTA_CATALOG,
+ LakeCatalog.TYPE_CUSTOM,
+ LakeFormats.DELTA_CATALOG,
"",
"",
confWithDocs(
@@ -154,7 +154,7 @@ public interface Docs {
private final String uri;
private final String confExtra;
- SparkCatalogTemplate(
+ LakeCatalogTemplate(
String displayName,
String description,
String catalogName,
@@ -195,14 +195,14 @@ public String getDescription() {
/** Labels for selection dialogs (same order as {@link #values()}). */
public static String[] displayNames() {
- return Arrays.stream(values()).map(SparkCatalogTemplate::getDisplayName).toArray(String[]::new);
+ return Arrays.stream(values()).map(LakeCatalogTemplate::getDisplayName).toArray(String[]::new);
}
- public static SparkCatalogTemplate fromDisplayName(String name) {
+ public static LakeCatalogTemplate fromDisplayName(String name) {
if (name == null) {
return null;
}
- for (SparkCatalogTemplate t : values()) {
+ for (LakeCatalogTemplate t : values()) {
if (t.displayName.equals(name)) {
return t;
}
@@ -214,7 +214,7 @@ public static SparkCatalogTemplate fromDisplayName(String name) {
* Apply this template onto {@code catalog}. Does not change the Hop metadata object name; leaves
* credential empty (never put secrets in presets).
*/
- public void applyTo(SparkCatalog catalog) {
+ public void applyTo(LakeCatalog catalog) {
Objects.requireNonNull(catalog, "catalog");
catalog.setCatalogName(catalogName);
catalog.setCatalogType(catalogType);
@@ -226,14 +226,14 @@ public void applyTo(SparkCatalog catalog) {
}
/**
- * True if the catalog looks customized relative to a fresh {@link SparkCatalog} default (used for
+ * True if the catalog looks customized relative to a fresh {@link LakeCatalog} default (used for
* overwrite confirmation).
*/
- public static boolean looksCustomized(SparkCatalog catalog) {
+ public static boolean looksCustomized(LakeCatalog catalog) {
if (catalog == null) {
return false;
}
- SparkCatalog def = new SparkCatalog();
+ LakeCatalog def = new LakeCatalog();
if (StringUtils.isNotEmpty(catalog.getCatalogName())) {
return true;
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInput.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInput.java
similarity index 84%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInput.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInput.java
index 741de475a1c..90ad3984e47 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInput.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInput.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.pipeline.Pipeline;
@@ -26,13 +26,12 @@
/**
* Metadata-only on Local engine. On the native Spark engine this becomes a lake table Dataset read.
*/
-public class SparkLakeTableInput
- extends BaseTransform {
+public class LakeTableInput extends BaseTransform {
- public SparkLakeTableInput(
+ public LakeTableInput(
TransformMeta transformMeta,
- SparkLakeTableInputMeta meta,
- SparkLakeTableInputData data,
+ LakeTableInputMeta meta,
+ LakeTableInputData data,
int copyNr,
PipelineMeta pipelineMeta,
Pipeline pipeline) {
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputData.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputData.java
similarity index 84%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputData.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputData.java
index f89ee00e6ed..58fa976d845 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputData.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputData.java
@@ -15,13 +15,13 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.pipeline.transform.BaseTransformData;
import org.apache.hop.pipeline.transform.ITransformData;
-public class SparkLakeTableInputData extends BaseTransformData implements ITransformData {
- public SparkLakeTableInputData() {
+public class LakeTableInputData extends BaseTransformData implements ITransformData {
+ public LakeTableInputData() {
super();
}
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputDialog.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputDialog.java
similarity index 91%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputDialog.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputDialog.java
index d6b2a1d65cf..7a5992e5add 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputDialog.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputDialog.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import java.util.ArrayList;
import java.util.List;
@@ -24,9 +24,9 @@
import org.apache.hop.core.util.Utils;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.i18n.BaseMessages;
+import org.apache.hop.lakehouse.LakeField;
+import org.apache.hop.lakehouse.LakeFormats;
import org.apache.hop.pipeline.PipelineMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.transforms.io.SparkField;
import org.apache.hop.ui.core.PropsUi;
import org.apache.hop.ui.core.dialog.BaseDialog;
import org.apache.hop.ui.core.widget.ColumnInfo;
@@ -46,10 +46,10 @@
import org.eclipse.swt.widgets.TableItem;
import org.eclipse.swt.widgets.Text;
-public class SparkLakeTableInputDialog extends BaseTransformDialog {
- private static final Class> PKG = SparkLakeTableInputMeta.class;
+public class LakeTableInputDialog extends BaseTransformDialog {
+ private static final Class> PKG = LakeTableInputMeta.class;
- private final SparkLakeTableInputMeta input;
+ private final LakeTableInputMeta input;
private CCombo wFormat;
private CCombo wIdentifierMode;
@@ -62,10 +62,10 @@ public class SparkLakeTableInputDialog extends BaseTransformDialog {
private Text wExtraOptions;
private TableView wFields;
- public SparkLakeTableInputDialog(
+ public LakeTableInputDialog(
Shell parent,
IVariables variables,
- SparkLakeTableInputMeta transformMeta,
+ LakeTableInputMeta transformMeta,
PipelineMeta pipelineMeta) {
super(parent, variables, transformMeta, pipelineMeta);
this.input = transformMeta;
@@ -110,13 +110,13 @@ public String open() {
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableInputDialog.Format", true);
wFormat = (CCombo) last;
- wFormat.setItems(new String[] {SparkLakeFormats.FORMAT_DELTA, SparkLakeFormats.FORMAT_ICEBERG});
+ wFormat.setItems(new String[] {LakeFormats.FORMAT_DELTA, LakeFormats.FORMAT_ICEBERG});
last =
labeledCombo(lsMod, middle, margin, last, "SparkLakeTableInputDialog.IdentifierMode", true);
wIdentifierMode = (CCombo) last;
wIdentifierMode.setItems(
- new String[] {SparkLakeTableInputMeta.MODE_PATH, SparkLakeTableInputMeta.MODE_TABLE});
+ new String[] {LakeTableInputMeta.MODE_PATH, LakeTableInputMeta.MODE_TABLE});
last = labeledTextVar(lsMod, middle, margin, last, "SparkLakeTableInputDialog.TablePath");
wTablePath = (TextVar) last;
@@ -134,9 +134,9 @@ public String open() {
wTimeTravelType = (CCombo) last;
wTimeTravelType.setItems(
new String[] {
- SparkLakeTableInputMeta.TIME_TRAVEL_NONE,
- SparkLakeTableInputMeta.TIME_TRAVEL_VERSION,
- SparkLakeTableInputMeta.TIME_TRAVEL_TIMESTAMP
+ LakeTableInputMeta.TIME_TRAVEL_NONE,
+ LakeTableInputMeta.TIME_TRAVEL_VERSION,
+ LakeTableInputMeta.TIME_TRAVEL_TIMESTAMP
});
last =
@@ -280,20 +280,19 @@ private TextVar labeledTextVar(
private void getData() {
wTransformName.setText(Const.NVL(transformName, ""));
- wFormat.setText(Const.NVL(input.getFormat(), SparkLakeFormats.FORMAT_DELTA));
- wIdentifierMode.setText(
- Const.NVL(input.getIdentifierMode(), SparkLakeTableInputMeta.MODE_PATH));
+ wFormat.setText(Const.NVL(input.getFormat(), LakeFormats.FORMAT_DELTA));
+ wIdentifierMode.setText(Const.NVL(input.getIdentifierMode(), LakeTableInputMeta.MODE_PATH));
wTablePath.setText(Const.NVL(input.getTablePath(), ""));
wTableIdentifier.setText(Const.NVL(input.getTableIdentifier(), ""));
wCatalogMetadataName.setText(Const.NVL(input.getCatalogMetadataName(), ""));
wTimeTravelType.setText(
- Const.NVL(input.getTimeTravelType(), SparkLakeTableInputMeta.TIME_TRAVEL_NONE));
+ Const.NVL(input.getTimeTravelType(), LakeTableInputMeta.TIME_TRAVEL_NONE));
wTimeTravelVersion.setText(Const.NVL(input.getTimeTravelVersion(), ""));
wTimeTravelTimestamp.setText(Const.NVL(input.getTimeTravelTimestamp(), ""));
wExtraOptions.setText(Const.NVL(input.getExtraOptions(), ""));
if (input.getFields() != null) {
int i = 0;
- for (SparkField f : input.getFields()) {
+ for (LakeField f : input.getFields()) {
TableItem item = wFields.table.getItem(i);
if (item == null) {
item = new TableItem(wFields.table, SWT.NONE);
@@ -331,10 +330,10 @@ private void ok() {
input.setTimeTravelVersion(wTimeTravelVersion.getText());
input.setTimeTravelTimestamp(wTimeTravelTimestamp.getText());
input.setExtraOptions(wExtraOptions.getText());
- List fields = new ArrayList<>();
+ List fields = new ArrayList<>();
for (int i = 0; i < wFields.nrNonEmpty(); i++) {
TableItem item = wFields.getNonEmpty(i);
- SparkField f = new SparkField();
+ LakeField f = new LakeField();
f.setName(item.getText(1));
f.setHopType(item.getText(2));
f.setLength(Const.toInt(item.getText(3), -1));
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMeta.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputMeta.java
similarity index 85%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMeta.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputMeta.java
index 78df4b7c1b2..7afc107068f 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMeta.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableInputMeta.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import java.util.ArrayList;
import java.util.List;
@@ -26,27 +26,26 @@
import org.apache.hop.core.exception.HopTransformException;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.lakehouse.LakeField;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.LakehouseConst;
import org.apache.hop.metadata.api.HopMetadataProperty;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.transform.BaseTransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.transforms.io.SparkField;
-import org.apache.hop.spark.util.SparkConst;
@Transform(
- id = SparkConst.SPARK_LAKE_TABLE_INPUT_PLUGIN_ID,
+ id = LakehouseConst.LAKE_TABLE_INPUT_PLUGIN_ID,
name = "i18n::SparkLakeTableInput.Name",
description = "i18n::SparkLakeTableInput.Description",
image = "spark-lake-table-input.svg",
categoryDescription = "i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.BigData",
keywords = "i18n::SparkLakeTableInput.Keyword",
documentationUrl = "/pipeline/transforms/spark-lake-table-input.html",
- supportedEngines = {SparkConst.PLUGIN_ID})
+ supportedEngines = {LakehouseConst.SPARK_ENGINE_ID})
@Getter
@Setter
-public class SparkLakeTableInputMeta
- extends BaseTransformMeta {
+public class LakeTableInputMeta extends BaseTransformMeta {
public static final String MODE_PATH = "PATH";
public static final String MODE_TABLE = "TABLE";
@@ -55,9 +54,9 @@ public class SparkLakeTableInputMeta
public static final String TIME_TRAVEL_VERSION = "VERSION";
public static final String TIME_TRAVEL_TIMESTAMP = "TIMESTAMP";
- /** {@link SparkLakeFormats#FORMAT_DELTA} or {@link SparkLakeFormats#FORMAT_ICEBERG} */
+ /** {@link LakeFormats#FORMAT_DELTA} or {@link LakeFormats#FORMAT_ICEBERG} */
@HopMetadataProperty(key = "format", injectionKey = "FORMAT")
- private String format = SparkLakeFormats.FORMAT_DELTA;
+ private String format = LakeFormats.FORMAT_DELTA;
/** PATH (v1 primary) or TABLE (catalog — later PRs) */
@HopMetadataProperty(key = "identifier_mode", injectionKey = "IDENTIFIER_MODE")
@@ -70,7 +69,7 @@ public class SparkLakeTableInputMeta
@HopMetadataProperty(key = "table_identifier", injectionKey = "TABLE_IDENTIFIER")
private String tableIdentifier;
- /** Hop SparkCatalog metadata name when mode is TABLE. */
+ /** Hop LakeCatalog metadata name when mode is TABLE. */
@HopMetadataProperty(key = "catalog_metadata_name", injectionKey = "CATALOG_METADATA_NAME")
private String catalogMetadataName;
@@ -101,9 +100,9 @@ public class SparkLakeTableInputMeta
key = "field",
injectionGroupKey = "FIELDS",
injectionGroupDescription = "SparkLakeTableInput.Injection.Group.Fields")
- private List fields = new ArrayList<>();
+ private List fields = new ArrayList<>();
- public SparkLakeTableInputMeta() {
+ public LakeTableInputMeta() {
super();
}
@@ -119,7 +118,7 @@ public boolean canStartWithoutInput() {
@Override
public String getDialogClassName() {
- return SparkLakeTableInputDialog.class.getName();
+ return LakeTableInputDialog.class.getName();
}
@Override
@@ -136,7 +135,7 @@ public void getFields(
return;
}
try {
- for (SparkField field : fields) {
+ for (LakeField field : fields) {
if (field.getName() != null && !field.getName().isEmpty()) {
inputRowMeta.addValueMeta(field.createValueMeta());
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenance.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenance.java
similarity index 83%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenance.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenance.java
index 894de666d32..6cc28e20bca 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenance.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenance.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.pipeline.Pipeline;
@@ -27,13 +27,13 @@
* Metadata-only on Local engine. On native Spark this runs OPTIMIZE / VACUUM / expire / DELETE as a
* zero-input action sink.
*/
-public class SparkLakeTableMaintenance
- extends BaseTransform {
+public class LakeTableMaintenance
+ extends BaseTransform {
- public SparkLakeTableMaintenance(
+ public LakeTableMaintenance(
TransformMeta transformMeta,
- SparkLakeTableMaintenanceMeta meta,
- SparkLakeTableMaintenanceData data,
+ LakeTableMaintenanceMeta meta,
+ LakeTableMaintenanceData data,
int copyNr,
PipelineMeta pipelineMeta,
Pipeline pipeline) {
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputData.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceData.java
similarity index 86%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputData.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceData.java
index f756b321288..adf5bc09a89 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputData.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceData.java
@@ -15,13 +15,13 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.pipeline.transform.BaseTransformData;
import org.apache.hop.pipeline.transform.ITransformData;
-public class SparkLakeTableOutputData extends BaseTransformData implements ITransformData {
- public SparkLakeTableOutputData() {
+public class LakeTableMaintenanceData extends BaseTransformData implements ITransformData {
+ public LakeTableMaintenanceData() {
super();
}
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceDialog.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceDialog.java
similarity index 88%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceDialog.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceDialog.java
index 8fce987c9d2..80d46154c64 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceDialog.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceDialog.java
@@ -15,15 +15,14 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.Const;
import org.apache.hop.core.util.Utils;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.i18n.BaseMessages;
+import org.apache.hop.lakehouse.LakeFormats;
import org.apache.hop.pipeline.PipelineMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.table.SparkMaintenanceSqlBuilder;
import org.apache.hop.ui.core.PropsUi;
import org.apache.hop.ui.core.dialog.BaseDialog;
import org.apache.hop.ui.core.widget.TextVar;
@@ -39,10 +38,10 @@
import org.eclipse.swt.widgets.Label;
import org.eclipse.swt.widgets.Shell;
-public class SparkLakeTableMaintenanceDialog extends BaseTransformDialog {
- private static final Class> PKG = SparkLakeTableMaintenanceMeta.class;
+public class LakeTableMaintenanceDialog extends BaseTransformDialog {
+ private static final Class> PKG = LakeTableMaintenanceMeta.class;
- private final SparkLakeTableMaintenanceMeta input;
+ private final LakeTableMaintenanceMeta input;
private CCombo wFormat;
private CCombo wIdentifierMode;
@@ -56,10 +55,10 @@ public class SparkLakeTableMaintenanceDialog extends BaseTransformDialog {
private TextVar wZOrderColumns;
private Button wAcknowledgeDestructive;
- public SparkLakeTableMaintenanceDialog(
+ public LakeTableMaintenanceDialog(
Shell parent,
IVariables variables,
- SparkLakeTableMaintenanceMeta transformMeta,
+ LakeTableMaintenanceMeta transformMeta,
PipelineMeta pipelineMeta) {
super(parent, variables, transformMeta, pipelineMeta);
this.input = transformMeta;
@@ -104,13 +103,13 @@ public String open() {
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableMaintenanceDialog.Format");
wFormat = (CCombo) last;
- wFormat.setItems(new String[] {SparkLakeFormats.FORMAT_DELTA, SparkLakeFormats.FORMAT_ICEBERG});
+ wFormat.setItems(new String[] {LakeFormats.FORMAT_DELTA, LakeFormats.FORMAT_ICEBERG});
last =
labeledCombo(lsMod, middle, margin, last, "SparkLakeTableMaintenanceDialog.IdentifierMode");
wIdentifierMode = (CCombo) last;
wIdentifierMode.setItems(
- new String[] {SparkLakeTableInputMeta.MODE_PATH, SparkLakeTableInputMeta.MODE_TABLE});
+ new String[] {LakeTableInputMeta.MODE_PATH, LakeTableInputMeta.MODE_TABLE});
last = labeledTextVar(lsMod, middle, margin, last, "SparkLakeTableMaintenanceDialog.TablePath");
wTablePath = (TextVar) last;
@@ -129,11 +128,11 @@ public String open() {
wOperation = (CCombo) last;
wOperation.setItems(
new String[] {
- SparkMaintenanceSqlBuilder.OP_OPTIMIZE,
- SparkMaintenanceSqlBuilder.OP_VACUUM,
- SparkMaintenanceSqlBuilder.OP_EXPIRE_SNAPSHOTS,
- SparkMaintenanceSqlBuilder.OP_REWRITE_MANIFESTS,
- SparkMaintenanceSqlBuilder.OP_DELETE_WHERE
+ LakeTableMaintenanceMeta.OP_OPTIMIZE,
+ LakeTableMaintenanceMeta.OP_VACUUM,
+ LakeTableMaintenanceMeta.OP_EXPIRE_SNAPSHOTS,
+ LakeTableMaintenanceMeta.OP_REWRITE_MANIFESTS,
+ LakeTableMaintenanceMeta.OP_DELETE_WHERE
});
last =
@@ -228,13 +227,12 @@ private TextVar labeledTextVar(
private void getData() {
wTransformName.setText(Const.NVL(transformName, ""));
- wFormat.setText(Const.NVL(input.getFormat(), SparkLakeFormats.FORMAT_DELTA));
- wIdentifierMode.setText(
- Const.NVL(input.getIdentifierMode(), SparkLakeTableInputMeta.MODE_PATH));
+ wFormat.setText(Const.NVL(input.getFormat(), LakeFormats.FORMAT_DELTA));
+ wIdentifierMode.setText(Const.NVL(input.getIdentifierMode(), LakeTableInputMeta.MODE_PATH));
wTablePath.setText(Const.NVL(input.getTablePath(), ""));
wTableIdentifier.setText(Const.NVL(input.getTableIdentifier(), ""));
wCatalogMetadataName.setText(Const.NVL(input.getCatalogMetadataName(), ""));
- wOperation.setText(Const.NVL(input.getOperation(), SparkMaintenanceSqlBuilder.OP_OPTIMIZE));
+ wOperation.setText(Const.NVL(input.getOperation(), LakeTableMaintenanceMeta.OP_OPTIMIZE));
wRetentionHours.setText(Const.NVL(input.getRetentionHours(), ""));
wRetainLast.setText(Const.NVL(input.getRetainLast(), "1"));
wWhereClause.setText(Const.NVL(input.getWhereClause(), ""));
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMeta.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMeta.java
similarity index 78%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMeta.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMeta.java
index e3632ec1977..98a5c35f384 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMeta.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMeta.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import lombok.Getter;
import lombok.Setter;
@@ -23,33 +23,38 @@
import org.apache.hop.core.exception.HopTransformException;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.LakehouseConst;
import org.apache.hop.metadata.api.HopMetadataProperty;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.transform.BaseTransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.table.SparkMaintenanceSqlBuilder;
-import org.apache.hop.spark.util.SparkConst;
@Transform(
- id = SparkConst.SPARK_LAKE_TABLE_MAINTENANCE_PLUGIN_ID,
+ id = LakehouseConst.LAKE_TABLE_MAINTENANCE_PLUGIN_ID,
name = "i18n::SparkLakeTableMaintenance.Name",
description = "i18n::SparkLakeTableMaintenance.Description",
image = "spark-lake-table-maintenance.svg",
categoryDescription = "i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.BigData",
keywords = "i18n::SparkLakeTableMaintenance.Keyword",
documentationUrl = "/pipeline/transforms/spark-lake-table-maintenance.html",
- supportedEngines = {SparkConst.PLUGIN_ID})
+ supportedEngines = {LakehouseConst.SPARK_ENGINE_ID})
@Getter
@Setter
-public class SparkLakeTableMaintenanceMeta
- extends BaseTransformMeta {
+public class LakeTableMaintenanceMeta
+ extends BaseTransformMeta {
+
+ public static final String OP_OPTIMIZE = "OPTIMIZE";
+ public static final String OP_VACUUM = "VACUUM";
+ public static final String OP_EXPIRE_SNAPSHOTS = "EXPIRE_SNAPSHOTS";
+ public static final String OP_REWRITE_MANIFESTS = "REWRITE_MANIFESTS";
+ public static final String OP_DELETE_WHERE = "DELETE_WHERE";
@HopMetadataProperty(key = "format", injectionKey = "FORMAT")
- private String format = SparkLakeFormats.FORMAT_DELTA;
+ private String format = LakeFormats.FORMAT_DELTA;
@HopMetadataProperty(key = "identifier_mode", injectionKey = "IDENTIFIER_MODE")
- private String identifierMode = SparkLakeTableInputMeta.MODE_PATH;
+ private String identifierMode = LakeTableInputMeta.MODE_PATH;
@HopMetadataProperty(key = "table_path", injectionKey = "TABLE_PATH")
private String tablePath;
@@ -60,9 +65,9 @@ public class SparkLakeTableMaintenanceMeta
@HopMetadataProperty(key = "catalog_metadata_name", injectionKey = "CATALOG_METADATA_NAME")
private String catalogMetadataName;
- /** {@link SparkMaintenanceSqlBuilder} operation constants. */
+ /** The OP_* operation constants above. */
@HopMetadataProperty(key = "operation", injectionKey = "OPERATION")
- private String operation = SparkMaintenanceSqlBuilder.OP_OPTIMIZE;
+ private String operation = LakeTableMaintenanceMeta.OP_OPTIMIZE;
/** Required for VACUUM / EXPIRE_SNAPSHOTS (hours). No silent default. */
@HopMetadataProperty(key = "retention_hours", injectionKey = "RETENTION_HOURS")
@@ -87,13 +92,13 @@ public class SparkLakeTableMaintenanceMeta
@HopMetadataProperty(key = "acknowledge_destructive", injectionKey = "ACKNOWLEDGE_DESTRUCTIVE")
private boolean acknowledgeDestructive;
- public SparkLakeTableMaintenanceMeta() {
+ public LakeTableMaintenanceMeta() {
super();
}
@Override
public String getDialogClassName() {
- return SparkLakeTableMaintenanceDialog.class.getName();
+ return LakeTableMaintenanceDialog.class.getName();
}
@Override
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMerge.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMerge.java
similarity index 84%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMerge.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMerge.java
index 5de744c18fb..bb022e9f137 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMerge.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMerge.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.pipeline.Pipeline;
@@ -27,13 +27,12 @@
* Metadata-only on Local engine. On native Spark this becomes a MERGE INTO action against a lake
* table.
*/
-public class SparkLakeTableMerge
- extends BaseTransform {
+public class LakeTableMerge extends BaseTransform {
- public SparkLakeTableMerge(
+ public LakeTableMerge(
TransformMeta transformMeta,
- SparkLakeTableMergeMeta meta,
- SparkLakeTableMergeData data,
+ LakeTableMergeMeta meta,
+ LakeTableMergeData data,
int copyNr,
PipelineMeta pipelineMeta,
Pipeline pipeline) {
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeData.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeData.java
similarity index 84%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeData.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeData.java
index f8c59640965..a59aeed45e6 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeData.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeData.java
@@ -15,13 +15,13 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.pipeline.transform.BaseTransformData;
import org.apache.hop.pipeline.transform.ITransformData;
-public class SparkLakeTableMergeData extends BaseTransformData implements ITransformData {
- public SparkLakeTableMergeData() {
+public class LakeTableMergeData extends BaseTransformData implements ITransformData {
+ public LakeTableMergeData() {
super();
}
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeDialog.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeDialog.java
similarity index 87%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeDialog.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeDialog.java
index 6768b4a98d0..c384e168fbd 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeDialog.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeDialog.java
@@ -15,15 +15,14 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.Const;
import org.apache.hop.core.util.Utils;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.i18n.BaseMessages;
+import org.apache.hop.lakehouse.LakeFormats;
import org.apache.hop.pipeline.PipelineMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.table.SparkMergeSqlBuilder;
import org.apache.hop.ui.core.PropsUi;
import org.apache.hop.ui.core.dialog.BaseDialog;
import org.apache.hop.ui.core.widget.TextVar;
@@ -40,10 +39,10 @@
import org.eclipse.swt.widgets.Shell;
import org.eclipse.swt.widgets.Text;
-public class SparkLakeTableMergeDialog extends BaseTransformDialog {
- private static final Class> PKG = SparkLakeTableMergeMeta.class;
+public class LakeTableMergeDialog extends BaseTransformDialog {
+ private static final Class> PKG = LakeTableMergeMeta.class;
- private final SparkLakeTableMergeMeta input;
+ private final LakeTableMergeMeta input;
private CCombo wFormat;
private CCombo wIdentifierMode;
@@ -56,10 +55,10 @@ public class SparkLakeTableMergeDialog extends BaseTransformDialog {
private CCombo wNotMatchedBySourceAction;
private Text wRawMergeSql;
- public SparkLakeTableMergeDialog(
+ public LakeTableMergeDialog(
Shell parent,
IVariables variables,
- SparkLakeTableMergeMeta transformMeta,
+ LakeTableMergeMeta transformMeta,
PipelineMeta pipelineMeta) {
super(parent, variables, transformMeta, pipelineMeta);
this.input = transformMeta;
@@ -104,12 +103,12 @@ public String open() {
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableMergeDialog.Format");
wFormat = (CCombo) last;
- wFormat.setItems(new String[] {SparkLakeFormats.FORMAT_DELTA, SparkLakeFormats.FORMAT_ICEBERG});
+ wFormat.setItems(new String[] {LakeFormats.FORMAT_DELTA, LakeFormats.FORMAT_ICEBERG});
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableMergeDialog.IdentifierMode");
wIdentifierMode = (CCombo) last;
wIdentifierMode.setItems(
- new String[] {SparkLakeTableInputMeta.MODE_PATH, SparkLakeTableInputMeta.MODE_TABLE});
+ new String[] {LakeTableInputMeta.MODE_PATH, LakeTableInputMeta.MODE_TABLE});
last = labeledTextVar(lsMod, middle, margin, last, "SparkLakeTableMergeDialog.TablePath");
wTablePath = (TextVar) last;
@@ -129,16 +128,16 @@ public String open() {
wMatchedAction = (CCombo) last;
wMatchedAction.setItems(
new String[] {
- SparkMergeSqlBuilder.MATCHED_UPDATE_ALL,
- SparkMergeSqlBuilder.MATCHED_DELETE,
- SparkMergeSqlBuilder.MATCHED_NONE
+ LakeTableMergeMeta.MATCHED_UPDATE_ALL,
+ LakeTableMergeMeta.MATCHED_DELETE,
+ LakeTableMergeMeta.MATCHED_NONE
});
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableMergeDialog.NotMatchedAction");
wNotMatchedAction = (CCombo) last;
wNotMatchedAction.setItems(
new String[] {
- SparkMergeSqlBuilder.NOT_MATCHED_INSERT_ALL, SparkMergeSqlBuilder.NOT_MATCHED_NONE
+ LakeTableMergeMeta.NOT_MATCHED_INSERT_ALL, LakeTableMergeMeta.NOT_MATCHED_NONE
});
last =
@@ -147,8 +146,8 @@ public String open() {
wNotMatchedBySourceAction = (CCombo) last;
wNotMatchedBySourceAction.setItems(
new String[] {
- SparkMergeSqlBuilder.NOT_MATCHED_BY_SOURCE_NONE,
- SparkMergeSqlBuilder.NOT_MATCHED_BY_SOURCE_DELETE
+ LakeTableMergeMeta.NOT_MATCHED_BY_SOURCE_NONE,
+ LakeTableMergeMeta.NOT_MATCHED_BY_SOURCE_DELETE
});
Label wlRaw = new Label(shell, SWT.RIGHT);
@@ -233,20 +232,19 @@ private TextVar labeledTextVar(
private void getData() {
wTransformName.setText(Const.NVL(transformName, ""));
- wFormat.setText(Const.NVL(input.getFormat(), SparkLakeFormats.FORMAT_DELTA));
- wIdentifierMode.setText(
- Const.NVL(input.getIdentifierMode(), SparkLakeTableInputMeta.MODE_PATH));
+ wFormat.setText(Const.NVL(input.getFormat(), LakeFormats.FORMAT_DELTA));
+ wIdentifierMode.setText(Const.NVL(input.getIdentifierMode(), LakeTableInputMeta.MODE_PATH));
wTablePath.setText(Const.NVL(input.getTablePath(), ""));
wTableIdentifier.setText(Const.NVL(input.getTableIdentifier(), ""));
wCatalogMetadataName.setText(Const.NVL(input.getCatalogMetadataName(), ""));
wMergeCondition.setText(Const.NVL(input.getMergeCondition(), "t.id = s.id"));
wMatchedAction.setText(
- Const.NVL(input.getMatchedAction(), SparkMergeSqlBuilder.MATCHED_UPDATE_ALL));
+ Const.NVL(input.getMatchedAction(), LakeTableMergeMeta.MATCHED_UPDATE_ALL));
wNotMatchedAction.setText(
- Const.NVL(input.getNotMatchedAction(), SparkMergeSqlBuilder.NOT_MATCHED_INSERT_ALL));
+ Const.NVL(input.getNotMatchedAction(), LakeTableMergeMeta.NOT_MATCHED_INSERT_ALL));
wNotMatchedBySourceAction.setText(
Const.NVL(
- input.getNotMatchedBySourceAction(), SparkMergeSqlBuilder.NOT_MATCHED_BY_SOURCE_NONE));
+ input.getNotMatchedBySourceAction(), LakeTableMergeMeta.NOT_MATCHED_BY_SOURCE_NONE));
wRawMergeSql.setText(Const.NVL(input.getRawMergeSql(), ""));
wTransformName.selectAll();
wTransformName.setFocus();
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMeta.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMeta.java
similarity index 69%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMeta.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMeta.java
index c2bbcd95da8..f790d21cd6c 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMeta.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMeta.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import lombok.Getter;
import lombok.Setter;
@@ -23,33 +23,41 @@
import org.apache.hop.core.exception.HopTransformException;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.LakehouseConst;
import org.apache.hop.metadata.api.HopMetadataProperty;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.transform.BaseTransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.table.SparkMergeSqlBuilder;
-import org.apache.hop.spark.util.SparkConst;
@Transform(
- id = SparkConst.SPARK_LAKE_TABLE_MERGE_PLUGIN_ID,
+ id = LakehouseConst.LAKE_TABLE_MERGE_PLUGIN_ID,
name = "i18n::SparkLakeTableMerge.Name",
description = "i18n::SparkLakeTableMerge.Description",
image = "spark-lake-table-merge.svg",
categoryDescription = "i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.BigData",
keywords = "i18n::SparkLakeTableMerge.Keyword",
documentationUrl = "/pipeline/transforms/spark-lake-table-merge.html",
- supportedEngines = {SparkConst.PLUGIN_ID})
+ supportedEngines = {LakehouseConst.SPARK_ENGINE_ID})
@Getter
@Setter
-public class SparkLakeTableMergeMeta
- extends BaseTransformMeta {
+public class LakeTableMergeMeta extends BaseTransformMeta {
+
+ public static final String MATCHED_UPDATE_ALL = "UPDATE_ALL";
+ public static final String MATCHED_DELETE = "DELETE";
+ public static final String MATCHED_NONE = "NONE";
+
+ public static final String NOT_MATCHED_INSERT_ALL = "INSERT_ALL";
+ public static final String NOT_MATCHED_NONE = "NONE";
+
+ public static final String NOT_MATCHED_BY_SOURCE_DELETE = "DELETE";
+ public static final String NOT_MATCHED_BY_SOURCE_NONE = "NONE";
@HopMetadataProperty(key = "format", injectionKey = "FORMAT")
- private String format = SparkLakeFormats.FORMAT_DELTA;
+ private String format = LakeFormats.FORMAT_DELTA;
@HopMetadataProperty(key = "identifier_mode", injectionKey = "IDENTIFIER_MODE")
- private String identifierMode = SparkLakeTableInputMeta.MODE_PATH;
+ private String identifierMode = LakeTableInputMeta.MODE_PATH;
@HopMetadataProperty(key = "table_path", injectionKey = "TABLE_PATH")
private String tablePath;
@@ -64,22 +72,21 @@ public class SparkLakeTableMergeMeta
@HopMetadataProperty(key = "merge_condition", injectionKey = "MERGE_CONDITION")
private String mergeCondition;
- /** {@link SparkMergeSqlBuilder#MATCHED_UPDATE_ALL}, DELETE, or NONE */
+ /** {@link #MATCHED_UPDATE_ALL}, DELETE, or NONE */
@HopMetadataProperty(key = "matched_action", injectionKey = "MATCHED_ACTION")
- private String matchedAction = SparkMergeSqlBuilder.MATCHED_UPDATE_ALL;
+ private String matchedAction = LakeTableMergeMeta.MATCHED_UPDATE_ALL;
- /** {@link SparkMergeSqlBuilder#NOT_MATCHED_INSERT_ALL} or NONE */
+ /** {@link #NOT_MATCHED_INSERT_ALL} or NONE */
@HopMetadataProperty(key = "not_matched_action", injectionKey = "NOT_MATCHED_ACTION")
- private String notMatchedAction = SparkMergeSqlBuilder.NOT_MATCHED_INSERT_ALL;
+ private String notMatchedAction = LakeTableMergeMeta.NOT_MATCHED_INSERT_ALL;
/**
- * Optional {@link SparkMergeSqlBuilder#NOT_MATCHED_BY_SOURCE_DELETE} (Delta; Iceberg support
- * varies). Default NONE.
+ * Optional {@link #NOT_MATCHED_BY_SOURCE_DELETE} (Delta; Iceberg support varies). Default NONE.
*/
@HopMetadataProperty(
key = "not_matched_by_source_action",
injectionKey = "NOT_MATCHED_BY_SOURCE_ACTION")
- private String notMatchedBySourceAction = SparkMergeSqlBuilder.NOT_MATCHED_BY_SOURCE_NONE;
+ private String notMatchedBySourceAction = LakeTableMergeMeta.NOT_MATCHED_BY_SOURCE_NONE;
/**
* Advanced: full MERGE SQL. When non-empty, overrides structured fields. Operator is trusted;
@@ -88,13 +95,13 @@ public class SparkLakeTableMergeMeta
@HopMetadataProperty(key = "raw_merge_sql", injectionKey = "RAW_MERGE_SQL")
private String rawMergeSql;
- public SparkLakeTableMergeMeta() {
+ public LakeTableMergeMeta() {
super();
}
@Override
public String getDialogClassName() {
- return SparkLakeTableMergeDialog.class.getName();
+ return LakeTableMergeDialog.class.getName();
}
@Override
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutput.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutput.java
similarity index 84%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutput.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutput.java
index d94c2b30947..b2802895e33 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutput.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutput.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.pipeline.Pipeline;
@@ -26,13 +26,12 @@
/**
* Metadata-only on Local engine. On the native Spark engine this becomes a lake table write action.
*/
-public class SparkLakeTableOutput
- extends BaseTransform {
+public class LakeTableOutput extends BaseTransform {
- public SparkLakeTableOutput(
+ public LakeTableOutput(
TransformMeta transformMeta,
- SparkLakeTableOutputMeta meta,
- SparkLakeTableOutputData data,
+ LakeTableOutputMeta meta,
+ LakeTableOutputData data,
int copyNr,
PipelineMeta pipelineMeta,
Pipeline pipeline) {
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceData.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputData.java
similarity index 83%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceData.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputData.java
index 35bcb120374..f5510ddb101 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceData.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputData.java
@@ -15,13 +15,13 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.pipeline.transform.BaseTransformData;
import org.apache.hop.pipeline.transform.ITransformData;
-public class SparkLakeTableMaintenanceData extends BaseTransformData implements ITransformData {
- public SparkLakeTableMaintenanceData() {
+public class LakeTableOutputData extends BaseTransformData implements ITransformData {
+ public LakeTableOutputData() {
super();
}
}
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputDialog.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputDialog.java
similarity index 89%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputDialog.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputDialog.java
index bbe24a6245a..dacfab15784 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputDialog.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputDialog.java
@@ -15,15 +15,14 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.Const;
import org.apache.hop.core.util.Utils;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.i18n.BaseMessages;
+import org.apache.hop.lakehouse.LakeFormats;
import org.apache.hop.pipeline.PipelineMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.transforms.io.SparkFileOutputMeta;
import org.apache.hop.ui.core.PropsUi;
import org.apache.hop.ui.core.dialog.BaseDialog;
import org.apache.hop.ui.core.widget.TextVar;
@@ -40,10 +39,10 @@
import org.eclipse.swt.widgets.Shell;
import org.eclipse.swt.widgets.Text;
-public class SparkLakeTableOutputDialog extends BaseTransformDialog {
- private static final Class> PKG = SparkLakeTableOutputMeta.class;
+public class LakeTableOutputDialog extends BaseTransformDialog {
+ private static final Class> PKG = LakeTableOutputMeta.class;
- private final SparkLakeTableOutputMeta input;
+ private final LakeTableOutputMeta input;
private CCombo wFormat;
private CCombo wIdentifierMode;
@@ -55,10 +54,10 @@ public class SparkLakeTableOutputDialog extends BaseTransformDialog {
private TextVar wCoalesce;
private Text wExtraOptions;
- public SparkLakeTableOutputDialog(
+ public LakeTableOutputDialog(
Shell parent,
IVariables variables,
- SparkLakeTableOutputMeta transformMeta,
+ LakeTableOutputMeta transformMeta,
PipelineMeta pipelineMeta) {
super(parent, variables, transformMeta, pipelineMeta);
this.input = transformMeta;
@@ -103,12 +102,12 @@ public String open() {
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableOutputDialog.Format");
wFormat = (CCombo) last;
- wFormat.setItems(new String[] {SparkLakeFormats.FORMAT_DELTA, SparkLakeFormats.FORMAT_ICEBERG});
+ wFormat.setItems(new String[] {LakeFormats.FORMAT_DELTA, LakeFormats.FORMAT_ICEBERG});
last = labeledCombo(lsMod, middle, margin, last, "SparkLakeTableOutputDialog.IdentifierMode");
wIdentifierMode = (CCombo) last;
wIdentifierMode.setItems(
- new String[] {SparkLakeTableInputMeta.MODE_PATH, SparkLakeTableInputMeta.MODE_TABLE});
+ new String[] {LakeTableInputMeta.MODE_PATH, LakeTableInputMeta.MODE_TABLE});
last = labeledTextVar(lsMod, middle, margin, last, "SparkLakeTableOutputDialog.TablePath");
wTablePath = (TextVar) last;
@@ -126,10 +125,10 @@ public String open() {
wSaveMode = (CCombo) last;
wSaveMode.setItems(
new String[] {
- SparkFileOutputMeta.MODE_ERROR,
- SparkFileOutputMeta.MODE_APPEND,
- SparkFileOutputMeta.MODE_OVERWRITE,
- SparkFileOutputMeta.MODE_IGNORE
+ LakeTableOutputMeta.MODE_ERROR,
+ LakeTableOutputMeta.MODE_APPEND,
+ LakeTableOutputMeta.MODE_OVERWRITE,
+ LakeTableOutputMeta.MODE_IGNORE
});
last = labeledTextVar(lsMod, middle, margin, last, "SparkLakeTableOutputDialog.PartitionBy");
@@ -219,13 +218,12 @@ private TextVar labeledTextVar(
private void getData() {
wTransformName.setText(Const.NVL(transformName, ""));
- wFormat.setText(Const.NVL(input.getFormat(), SparkLakeFormats.FORMAT_DELTA));
- wIdentifierMode.setText(
- Const.NVL(input.getIdentifierMode(), SparkLakeTableInputMeta.MODE_PATH));
+ wFormat.setText(Const.NVL(input.getFormat(), LakeFormats.FORMAT_DELTA));
+ wIdentifierMode.setText(Const.NVL(input.getIdentifierMode(), LakeTableInputMeta.MODE_PATH));
wTablePath.setText(Const.NVL(input.getTablePath(), ""));
wTableIdentifier.setText(Const.NVL(input.getTableIdentifier(), ""));
wCatalogMetadataName.setText(Const.NVL(input.getCatalogMetadataName(), ""));
- wSaveMode.setText(Const.NVL(input.getSaveMode(), SparkFileOutputMeta.MODE_ERROR));
+ wSaveMode.setText(Const.NVL(input.getSaveMode(), LakeTableOutputMeta.MODE_ERROR));
wPartitionBy.setText(Const.NVL(input.getPartitionByColumns(), ""));
wCoalesce.setText(Const.NVL(input.getCoalescePartitions(), ""));
wExtraOptions.setText(Const.NVL(input.getExtraOptions(), ""));
diff --git a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMeta.java b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMeta.java
similarity index 73%
rename from plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMeta.java
rename to plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMeta.java
index 70586b5c970..02fd1abf6e5 100644
--- a/plugins/engines/spark/src/main/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMeta.java
+++ b/plugins/tech/lakehouse/src/main/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMeta.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import lombok.Getter;
import lombok.Setter;
@@ -23,34 +23,39 @@
import org.apache.hop.core.exception.HopTransformException;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.LakehouseConst;
import org.apache.hop.metadata.api.HopMetadataProperty;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.transform.BaseTransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
-import org.apache.hop.spark.table.SparkLakeFormats;
-import org.apache.hop.spark.transforms.io.SparkFileOutputMeta;
-import org.apache.hop.spark.util.SparkConst;
@Transform(
- id = SparkConst.SPARK_LAKE_TABLE_OUTPUT_PLUGIN_ID,
+ id = LakehouseConst.LAKE_TABLE_OUTPUT_PLUGIN_ID,
name = "i18n::SparkLakeTableOutput.Name",
description = "i18n::SparkLakeTableOutput.Description",
image = "spark-lake-table-output.svg",
categoryDescription = "i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.BigData",
keywords = "i18n::SparkLakeTableOutput.Keyword",
documentationUrl = "/pipeline/transforms/spark-lake-table-output.html",
- supportedEngines = {SparkConst.PLUGIN_ID})
+ supportedEngines = {LakehouseConst.SPARK_ENGINE_ID})
@Getter
@Setter
-public class SparkLakeTableOutputMeta
- extends BaseTransformMeta {
+public class LakeTableOutputMeta extends BaseTransformMeta {
- /** {@link SparkLakeFormats#FORMAT_DELTA} or {@link SparkLakeFormats#FORMAT_ICEBERG} */
+ /** Save modes, using the names of Spark's {@code SaveMode}. */
+ public static final String MODE_OVERWRITE = "Overwrite";
+
+ public static final String MODE_APPEND = "Append";
+ public static final String MODE_IGNORE = "Ignore";
+ public static final String MODE_ERROR = "ErrorIfExists";
+
+ /** {@link LakeFormats#FORMAT_DELTA} or {@link LakeFormats#FORMAT_ICEBERG} */
@HopMetadataProperty(key = "format", injectionKey = "FORMAT")
- private String format = SparkLakeFormats.FORMAT_DELTA;
+ private String format = LakeFormats.FORMAT_DELTA;
@HopMetadataProperty(key = "identifier_mode", injectionKey = "IDENTIFIER_MODE")
- private String identifierMode = SparkLakeTableInputMeta.MODE_PATH;
+ private String identifierMode = LakeTableInputMeta.MODE_PATH;
@HopMetadataProperty(key = "table_path", injectionKey = "TABLE_PATH")
private String tablePath;
@@ -62,11 +67,11 @@ public class SparkLakeTableOutputMeta
private String catalogMetadataName;
/**
- * Default {@link SparkFileOutputMeta#MODE_ERROR} (ErrorIfExists) — ACID tables must not default
- * to destructive overwrite.
+ * Default {@link #MODE_ERROR} (ErrorIfExists) — ACID tables must not default to destructive
+ * overwrite.
*/
@HopMetadataProperty(key = "save_mode", injectionKey = "SAVE_MODE")
- private String saveMode = SparkFileOutputMeta.MODE_ERROR;
+ private String saveMode = MODE_ERROR;
@HopMetadataProperty(key = "partition_by", injectionKey = "PARTITION_BY")
private String partitionByColumns;
@@ -78,13 +83,13 @@ public class SparkLakeTableOutputMeta
@HopMetadataProperty(key = "extra_options", injectionKey = "EXTRA_OPTIONS")
private String extraOptions;
- public SparkLakeTableOutputMeta() {
+ public LakeTableOutputMeta() {
super();
}
@Override
public String getDialogClassName() {
- return SparkLakeTableOutputDialog.class.getName();
+ return LakeTableOutputDialog.class.getName();
}
@Override
diff --git a/plugins/engines/spark/src/main/resources/org/apache/hop/spark/metadata/messages/messages_en_US.properties b/plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/metadata/messages/messages_en_US.properties
similarity index 100%
rename from plugins/engines/spark/src/main/resources/org/apache/hop/spark/metadata/messages/messages_en_US.properties
rename to plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/metadata/messages/messages_en_US.properties
diff --git a/plugins/engines/spark/src/main/resources/org/apache/hop/spark/metadata/messages/messages_pt_BR.properties b/plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/metadata/messages/messages_pt_BR.properties
similarity index 100%
rename from plugins/engines/spark/src/main/resources/org/apache/hop/spark/metadata/messages/messages_pt_BR.properties
rename to plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/metadata/messages/messages_pt_BR.properties
diff --git a/plugins/engines/spark/src/main/resources/org/apache/hop/spark/transforms/table/messages/messages_en_US.properties b/plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/transforms/messages/messages_en_US.properties
similarity index 100%
rename from plugins/engines/spark/src/main/resources/org/apache/hop/spark/transforms/table/messages/messages_en_US.properties
rename to plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/transforms/messages/messages_en_US.properties
diff --git a/plugins/engines/spark/src/main/resources/org/apache/hop/spark/transforms/table/messages/messages_pt_BR.properties b/plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/transforms/messages/messages_pt_BR.properties
similarity index 100%
rename from plugins/engines/spark/src/main/resources/org/apache/hop/spark/transforms/table/messages/messages_pt_BR.properties
rename to plugins/tech/lakehouse/src/main/resources/org/apache/hop/lakehouse/transforms/messages/messages_pt_BR.properties
diff --git a/plugins/engines/spark/src/main/resources/spark-catalog.svg b/plugins/tech/lakehouse/src/main/resources/spark-catalog.svg
similarity index 100%
rename from plugins/engines/spark/src/main/resources/spark-catalog.svg
rename to plugins/tech/lakehouse/src/main/resources/spark-catalog.svg
diff --git a/plugins/engines/spark/src/main/resources/spark-lake-table-input.svg b/plugins/tech/lakehouse/src/main/resources/spark-lake-table-input.svg
similarity index 100%
rename from plugins/engines/spark/src/main/resources/spark-lake-table-input.svg
rename to plugins/tech/lakehouse/src/main/resources/spark-lake-table-input.svg
diff --git a/plugins/engines/spark/src/main/resources/spark-lake-table-maintenance.svg b/plugins/tech/lakehouse/src/main/resources/spark-lake-table-maintenance.svg
similarity index 100%
rename from plugins/engines/spark/src/main/resources/spark-lake-table-maintenance.svg
rename to plugins/tech/lakehouse/src/main/resources/spark-lake-table-maintenance.svg
diff --git a/plugins/engines/spark/src/main/resources/spark-lake-table-merge.svg b/plugins/tech/lakehouse/src/main/resources/spark-lake-table-merge.svg
similarity index 100%
rename from plugins/engines/spark/src/main/resources/spark-lake-table-merge.svg
rename to plugins/tech/lakehouse/src/main/resources/spark-lake-table-merge.svg
diff --git a/plugins/engines/spark/src/main/resources/spark-lake-table-output.svg b/plugins/tech/lakehouse/src/main/resources/spark-lake-table-output.svg
similarity index 100%
rename from plugins/engines/spark/src/main/resources/spark-lake-table-output.svg
rename to plugins/tech/lakehouse/src/main/resources/spark-lake-table-output.svg
diff --git a/plugins/tech/lakehouse/src/main/resources/version.xml b/plugins/tech/lakehouse/src/main/resources/version.xml
new file mode 100644
index 00000000000..ee1c2377fe3
--- /dev/null
+++ b/plugins/tech/lakehouse/src/main/resources/version.xml
@@ -0,0 +1,19 @@
+
+
+
+${project.version}
\ No newline at end of file
diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplateTest.java b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplateTest.java
similarity index 53%
rename from plugins/engines/spark/src/test/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplateTest.java
rename to plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplateTest.java
index 604e87e4d43..66fb87558b5 100644
--- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/metadata/template/SparkCatalogTemplateTest.java
+++ b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/metadata/template/LakeCatalogTemplateTest.java
@@ -15,26 +15,26 @@
* limitations under the License.
*/
-package org.apache.hop.spark.metadata.template;
+package org.apache.hop.lakehouse.metadata.template;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
-import org.apache.hop.spark.metadata.SparkCatalog;
-import org.apache.hop.spark.table.SparkLakeFormats;
+import org.apache.hop.lakehouse.LakeFormats;
+import org.apache.hop.lakehouse.metadata.LakeCatalog;
import org.junit.jupiter.api.Test;
-class SparkCatalogTemplateTest {
+class LakeCatalogTemplateTest {
@Test
void icebergHadoopLocalSetsWarehouseAndType() {
- SparkCatalog cat = new SparkCatalog();
- SparkCatalogTemplate.ICEBERG_HADOOP_LOCAL.applyTo(cat);
+ LakeCatalog cat = new LakeCatalog();
+ LakeCatalogTemplate.ICEBERG_HADOOP_LOCAL.applyTo(cat);
assertEquals("lake", cat.getCatalogName());
- assertEquals(SparkCatalog.TYPE_HADOOP, cat.getCatalogType());
- assertEquals(SparkLakeFormats.ICEBERG_CATALOG, cat.getImplementation());
+ assertEquals(LakeCatalog.TYPE_HADOOP, cat.getCatalogType());
+ assertEquals(LakeFormats.ICEBERG_CATALOG, cat.getImplementation());
assertEquals("file:///tmp/hop-warehouse", cat.getWarehouse());
assertEquals("", cat.getUri());
assertEquals("", cat.getCredential());
@@ -42,71 +42,71 @@ void icebergHadoopLocalSetsWarehouseAndType() {
@Test
void icebergRestSetsUri() {
- SparkCatalog cat = new SparkCatalog();
- SparkCatalogTemplate.ICEBERG_REST.applyTo(cat);
- assertEquals(SparkCatalog.TYPE_REST, cat.getCatalogType());
+ LakeCatalog cat = new LakeCatalog();
+ LakeCatalogTemplate.ICEBERG_REST.applyTo(cat);
+ assertEquals(LakeCatalog.TYPE_REST, cat.getCatalogType());
assertEquals("https://catalog.example.com/v1", cat.getUri());
assertTrue(cat.getWarehouse() == null || cat.getWarehouse().isEmpty());
}
@Test
void icebergRestAuthLeavesCredentialEmptyAndDocumentsToken() {
- SparkCatalog cat = new SparkCatalog();
+ LakeCatalog cat = new LakeCatalog();
cat.setCredential("should-be-cleared");
- SparkCatalogTemplate.ICEBERG_REST_AUTH.applyTo(cat);
- assertEquals(SparkCatalog.TYPE_REST, cat.getCatalogType());
+ LakeCatalogTemplate.ICEBERG_REST_AUTH.applyTo(cat);
+ assertEquals(LakeCatalog.TYPE_REST, cat.getCatalogType());
assertEquals("", cat.getCredential());
assertTrue(cat.getConfExtra().contains("Credential"));
}
@Test
void objectStoreSetsS3aAndIoImpl() {
- SparkCatalog cat = new SparkCatalog();
- SparkCatalogTemplate.ICEBERG_HADOOP_OBJECT_STORE.applyTo(cat);
+ LakeCatalog cat = new LakeCatalog();
+ LakeCatalogTemplate.ICEBERG_HADOOP_OBJECT_STORE.applyTo(cat);
assertEquals("s3a://bucket/warehouse", cat.getWarehouse());
assertTrue(cat.getConfExtra().contains("io-impl="));
}
@Test
void hiveAndGlueAreAdvancedTypes() {
- SparkCatalog hive = new SparkCatalog();
- SparkCatalogTemplate.HIVE_METASTORE.applyTo(hive);
- assertEquals(SparkCatalog.TYPE_HIVE, hive.getCatalogType());
+ LakeCatalog hive = new LakeCatalog();
+ LakeCatalogTemplate.HIVE_METASTORE.applyTo(hive);
+ assertEquals(LakeCatalog.TYPE_HIVE, hive.getCatalogType());
assertEquals("hive", hive.getCatalogName());
assertTrue(hive.getConfExtra().contains("thrift://"));
assertTrue(hive.getConfExtra().startsWith("# docs:"));
- SparkCatalog glue = new SparkCatalog();
- SparkCatalogTemplate.AWS_GLUE.applyTo(glue);
- assertEquals(SparkCatalog.TYPE_GLUE, glue.getCatalogType());
+ LakeCatalog glue = new LakeCatalog();
+ LakeCatalogTemplate.AWS_GLUE.applyTo(glue);
+ assertEquals(LakeCatalog.TYPE_GLUE, glue.getCatalogType());
assertEquals("glue", glue.getCatalogName());
assertTrue(glue.getConfExtra().contains("# docs:"));
}
@Test
void nessieUnityAndDeltaAdvancedTemplates() {
- SparkCatalog nessie = new SparkCatalog();
- SparkCatalogTemplate.NESSIE.applyTo(nessie);
- assertEquals(SparkCatalog.TYPE_CUSTOM, nessie.getCatalogType());
+ LakeCatalog nessie = new LakeCatalog();
+ LakeCatalogTemplate.NESSIE.applyTo(nessie);
+ assertEquals(LakeCatalog.TYPE_CUSTOM, nessie.getCatalogType());
assertEquals("nessie", nessie.getCatalogName());
assertTrue(nessie.getConfExtra().contains("NessieCatalog"));
- assertTrue(nessie.getConfExtra().contains(SparkCatalogTemplate.Docs.NESSIE_SPARK));
+ assertTrue(nessie.getConfExtra().contains(LakeCatalogTemplate.Docs.NESSIE_SPARK));
- SparkCatalog unity = new SparkCatalog();
- SparkCatalogTemplate.DATABRICKS_UNITY.applyTo(unity);
+ LakeCatalog unity = new LakeCatalog();
+ LakeCatalogTemplate.DATABRICKS_UNITY.applyTo(unity);
assertEquals("unity", unity.getCatalogName());
- assertTrue(unity.getConfExtra().contains(SparkCatalogTemplate.Docs.DATABRICKS_UNITY));
+ assertTrue(unity.getConfExtra().contains(LakeCatalogTemplate.Docs.DATABRICKS_UNITY));
- SparkCatalog delta = new SparkCatalog();
- SparkCatalogTemplate.DELTA_NAMED_CATALOG.applyTo(delta);
- assertEquals(SparkLakeFormats.DELTA_CATALOG, delta.getImplementation());
- assertTrue(delta.getConfExtra().contains(SparkCatalogTemplate.Docs.DELTA));
+ LakeCatalog delta = new LakeCatalog();
+ LakeCatalogTemplate.DELTA_NAMED_CATALOG.applyTo(delta);
+ assertEquals(LakeFormats.DELTA_CATALOG, delta.getImplementation());
+ assertTrue(delta.getConfExtra().contains(LakeCatalogTemplate.Docs.DELTA));
}
@Test
void everyTemplateIncludesDocsCommentInConfExtra() {
- for (SparkCatalogTemplate t : SparkCatalogTemplate.values()) {
- SparkCatalog cat = new SparkCatalog();
+ for (LakeCatalogTemplate t : LakeCatalogTemplate.values()) {
+ LakeCatalog cat = new LakeCatalog();
t.applyTo(cat);
assertTrue(
cat.getConfExtra() != null && cat.getConfExtra().contains("# docs:"),
@@ -117,39 +117,38 @@ void everyTemplateIncludesDocsCommentInConfExtra() {
@Test
void confWithDocsFormatsHeaderAndBody() {
assertEquals(
- "# docs: https://example.com",
- SparkCatalogTemplate.confWithDocs("https://example.com", ""));
+ "# docs: https://example.com", LakeCatalogTemplate.confWithDocs("https://example.com", ""));
assertEquals(
"# docs: https://example.com\nuri=thrift://x",
- SparkCatalogTemplate.confWithDocs("https://example.com", "uri=thrift://x"));
+ LakeCatalogTemplate.confWithDocs("https://example.com", "uri=thrift://x"));
}
@Test
void fromDisplayNameRoundTrip() {
- for (SparkCatalogTemplate t : SparkCatalogTemplate.values()) {
- assertEquals(t, SparkCatalogTemplate.fromDisplayName(t.getDisplayName()));
+ for (LakeCatalogTemplate t : LakeCatalogTemplate.values()) {
+ assertEquals(t, LakeCatalogTemplate.fromDisplayName(t.getDisplayName()));
}
- assertEquals(SparkCatalogTemplate.values().length, SparkCatalogTemplate.displayNames().length);
+ assertEquals(LakeCatalogTemplate.values().length, LakeCatalogTemplate.displayNames().length);
}
@Test
void looksCustomizedDetectsEdits() {
- SparkCatalog fresh = new SparkCatalog();
- assertFalse(SparkCatalogTemplate.looksCustomized(fresh));
+ LakeCatalog fresh = new LakeCatalog();
+ assertFalse(LakeCatalogTemplate.looksCustomized(fresh));
- SparkCatalog withName = new SparkCatalog();
+ LakeCatalog withName = new LakeCatalog();
withName.setCatalogName("lake");
- assertTrue(SparkCatalogTemplate.looksCustomized(withName));
+ assertTrue(LakeCatalogTemplate.looksCustomized(withName));
- SparkCatalog withWarehouse = new SparkCatalog();
+ LakeCatalog withWarehouse = new LakeCatalog();
withWarehouse.setWarehouse("file:///tmp/wh");
- assertTrue(SparkCatalogTemplate.looksCustomized(withWarehouse));
+ assertTrue(LakeCatalogTemplate.looksCustomized(withWarehouse));
}
@Test
void everyTemplateAppliesWithoutNpe() {
- for (SparkCatalogTemplate t : SparkCatalogTemplate.values()) {
- SparkCatalog cat = new SparkCatalog();
+ for (LakeCatalogTemplate t : LakeCatalogTemplate.values()) {
+ LakeCatalog cat = new LakeCatalog();
t.applyTo(cat);
assertNotNull(cat.getCatalogType());
assertNotNull(cat.getImplementation());
diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMetaInjectionTest.java b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableInputMetaInjectionTest.java
similarity index 86%
rename from plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMetaInjectionTest.java
rename to plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableInputMetaInjectionTest.java
index 96751c07e23..fbd16f28e94 100644
--- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableInputMetaInjectionTest.java
+++ b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableInputMetaInjectionTest.java
@@ -15,28 +15,27 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import java.util.ArrayList;
import java.util.List;
import org.apache.hop.core.injection.BaseMetadataInjectionTestJunit5;
import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
-import org.apache.hop.spark.transforms.io.SparkField;
+import org.apache.hop.lakehouse.LakeField;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
-class SparkLakeTableInputMetaInjectionTest
- extends BaseMetadataInjectionTestJunit5 {
+class LakeTableInputMetaInjectionTest extends BaseMetadataInjectionTestJunit5 {
@RegisterExtension
static RestoreHopEngineEnvironmentExtension env = new RestoreHopEngineEnvironmentExtension();
@BeforeEach
void setup() throws Exception {
- SparkLakeTableInputMeta meta = new SparkLakeTableInputMeta();
- List fields = new ArrayList<>();
- fields.add(new SparkField());
+ LakeTableInputMeta meta = new LakeTableInputMeta();
+ List fields = new ArrayList<>();
+ fields.add(new LakeField());
meta.setFields(fields);
setup(meta);
}
diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMetaInjectionTest.java b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMetaInjectionTest.java
similarity index 89%
rename from plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMetaInjectionTest.java
rename to plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMetaInjectionTest.java
index abb2adcf80b..ed020b3a0d8 100644
--- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMaintenanceMetaInjectionTest.java
+++ b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMaintenanceMetaInjectionTest.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.injection.BaseMetadataInjectionTestJunit5;
import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
@@ -23,15 +23,15 @@
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
-class SparkLakeTableMaintenanceMetaInjectionTest
- extends BaseMetadataInjectionTestJunit5 {
+class LakeTableMaintenanceMetaInjectionTest
+ extends BaseMetadataInjectionTestJunit5 {
@RegisterExtension
static RestoreHopEngineEnvironmentExtension env = new RestoreHopEngineEnvironmentExtension();
@BeforeEach
void setup() throws Exception {
- setup(new SparkLakeTableMaintenanceMeta());
+ setup(new LakeTableMaintenanceMeta());
}
@Test
diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMetaInjectionTest.java b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMetaInjectionTest.java
similarity index 90%
rename from plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMetaInjectionTest.java
rename to plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMetaInjectionTest.java
index 23183d79795..c8cc89383c3 100644
--- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableMergeMetaInjectionTest.java
+++ b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableMergeMetaInjectionTest.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.injection.BaseMetadataInjectionTestJunit5;
import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
@@ -23,15 +23,14 @@
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
-class SparkLakeTableMergeMetaInjectionTest
- extends BaseMetadataInjectionTestJunit5 {
+class LakeTableMergeMetaInjectionTest extends BaseMetadataInjectionTestJunit5 {
@RegisterExtension
static RestoreHopEngineEnvironmentExtension env = new RestoreHopEngineEnvironmentExtension();
@BeforeEach
void setup() throws Exception {
- setup(new SparkLakeTableMergeMeta());
+ setup(new LakeTableMergeMeta());
}
@Test
diff --git a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMetaInjectionTest.java b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMetaInjectionTest.java
similarity index 89%
rename from plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMetaInjectionTest.java
rename to plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMetaInjectionTest.java
index 2dc9f607409..b00472635fc 100644
--- a/plugins/engines/spark/src/test/java/org/apache/hop/spark/transforms/table/SparkLakeTableOutputMetaInjectionTest.java
+++ b/plugins/tech/lakehouse/src/test/java/org/apache/hop/lakehouse/transforms/LakeTableOutputMetaInjectionTest.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.hop.spark.transforms.table;
+package org.apache.hop.lakehouse.transforms;
import org.apache.hop.core.injection.BaseMetadataInjectionTestJunit5;
import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
@@ -23,15 +23,15 @@
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
-class SparkLakeTableOutputMetaInjectionTest
- extends BaseMetadataInjectionTestJunit5 {
+class LakeTableOutputMetaInjectionTest
+ extends BaseMetadataInjectionTestJunit5 {
@RegisterExtension
static RestoreHopEngineEnvironmentExtension env = new RestoreHopEngineEnvironmentExtension();
@BeforeEach
void setup() throws Exception {
- setup(new SparkLakeTableOutputMeta());
+ setup(new LakeTableOutputMeta());
}
@Test
diff --git a/plugins/tech/pom.xml b/plugins/tech/pom.xml
index 08deec71be7..0d53916f8e4 100644
--- a/plugins/tech/pom.xml
+++ b/plugins/tech/pom.xml
@@ -43,6 +43,7 @@
git-vfs
google
hadoop
+ lakehouse
minio
mongodb
neo4j