This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new 620d7c486f spark performance improvements, fixes #8369 (#8370)
620d7c486f is described below
commit 620d7c486fa81158b07222de9b4bb71457546bdc
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Mon Sep 14 18:03:20 2026 +0200
spark performance improvements, fixes #8369 (#8370)
* spark performance improvements, fixes #8369
* some more performance tweaks
---
.../spark/getting-started-with-native-spark.adoc | 2 +
.../ROOT/pages/pipeline/transforms/memgroupby.adoc | 2 +-
.../ROOT/pages/pipeline/transforms/mergejoin.adoc | 2 +-
.../ROOT/pages/pipeline/transforms/sort.adoc | 2 +-
.../transforms/spark-lake-table-input.adoc | 2 +-
.../transforms/spark-lake-table-maintenance.adoc | 2 +-
.../transforms/spark-lake-table-merge.adoc | 2 +-
.../transforms/spark-lake-table-output.adoc | 2 +-
.../ROOT/pages/pipeline/transforms/spark-sql.adoc | 2 +-
.../ROOT/pages/pipeline/transforms/uniquerows.adoc | 2 +-
.../hop/pipeline/transform/BaseTransform.java | 9 +-
.../apache/hop/spark/core/HopMapPartitionsFn.java | 449 +++++++++++++--------
.../hop/spark/core/HopSparkRowConverter.java | 91 +++++
.../hop/spark/core/SparkDirectRowHandler.java | 101 +++++
.../apache/hop/spark/core/SparkNativeMetrics.java | 178 +-------
.../hop/spark/core/SparkNativeMetricsListener.java | 350 ++++++++++++++++
.../hop/spark/engines/SparkPipelineEngine.java | 15 +
.../pipeline/HopPipelineMetaToSparkConverter.java | 10 +
.../handler/SparkBaseTransformHandler.java | 8 +-
.../pipeline/handler/SparkFileInputHandler.java | 203 +++++++++-
.../hop/spark/core/SparkNativeMetricsTest.java | 155 +++++--
.../pipeline/handler/SparkFileIoHandlersTest.java | 8 +
22 files changed, 1202 insertions(+), 395 deletions(-)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/getting-started-with-native-spark.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/getting-started-with-native-spark.adoc
index f79716f328..219729e64a 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/getting-started-with-native-spark.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/spark/getting-started-with-native-spark.adoc
@@ -105,6 +105,8 @@ If you need streaming, windowing, or Beam-only connectors,
stay on the Beam engi
** applies a *native handler* that rewrites the hop as Spark Dataset
operations, or
** wraps the transform in `mapPartitions` so each Spark partition runs that
transform in a tiny local Hop pipeline.
. File outputs and other Spark *actions* materialise the graph (read →
transform → write).
+
+Each transform's reference page says which of the two it gets: a double check
on the *Native Spark* tag means the engine has its own implementation
(xref:pipeline/transforms/sort.adoc[Sort rows],
xref:pipeline/transforms/mergejoin.adoc[Merge join],
xref:pipeline/transforms/memgroupby.adoc[Memory group by],
xref:pipeline/transforms/uniquerows.adoc[Unique rows], the Spark file and lake
table transforms and xref:pipeline/transforms/spark-sql.adoc[Spark SQL]); a
single check means the Hop tr [...]
. Transform metrics are collected via Spark accumulators and shown in Hop like
any other engine.
You still design pipelines *visually* in Hop. Spark’s role is the distributed
execution engine and the Dataset runtime under the hood.
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
index a0eece82ad..4ae7c9c8d6 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: This transform allows you to do aggregations in memory such as
sum,max,min,...
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=yes, Beam
Spark=yes, Beam Flink=yes, Beam Dataflow=yes
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=native, Beam
Spark=yes, Beam Flink=yes, Beam Dataflow=yes
= image:transforms/icons/memorygroupby.svg[Memory Group By transform Icon,
role="image-doc-icon"] Memory Group By
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/mergejoin.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/mergejoin.adoc
index d6f671b6a9..71ea11db90 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/mergejoin.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/mergejoin.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: The Merge Join transform performs a classic merge join between
data sets with data coming from two different input transforms.
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=yes, Beam
Spark=yes, Beam Flink=yes, Beam Dataflow=yes
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=native, Beam
Spark=yes, Beam Flink=yes, Beam Dataflow=yes
= image:transforms/icons/mergejoin.svg[Merge Join transform Icon,
role="image-doc-icon"] Merge Join
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
index e7e42ad967..aecbadd6d0 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: The Sort Rows transform sorts rows based on the fields you
specify and on whether they should be sorted in ascending or descending order.
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=yes, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/sortrows.svg[Sort Rows transform Icon,
role="image-doc-icon"] Sort Rows
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc
index 037527b49f..89813b1588 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-input.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: Spark lake table input reads Delta Lake or Apache Iceberg tables
on the native Spark pipeline engine, with optional time travel and projection.
-:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=yes, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/spark-lake-table-input.svg[Spark lake table input
transform Icon, role="image-doc-icon"] Spark lake table input
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc
index 5705e0baa6..cefaee6870 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-maintenance.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: Spark lake table maintenance runs OPTIMIZE, VACUUM, expire
snapshots, and related operations on Delta Lake or Apache Iceberg tables.
-:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=yes, Beam
Spark=yes, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/spark-lake-table-maintenance.svg[Spark lake table
maintenance transform Icon, role="image-doc-icon"] Spark lake table maintenance
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc
index 55718d33dd..53704ae1f7 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-merge.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: Spark lake table merge runs MERGE INTO against Delta Lake or
Apache Iceberg tables on the native Spark pipeline engine.
-:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=yes, Beam
Spark=yes, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/spark-lake-table-merge.svg[Spark lake table merge
transform Icon, role="image-doc-icon"] Spark lake table merge
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc
index 6c1ab554c1..6aeaf6a910 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-lake-table-output.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: Spark lake table output writes Delta Lake or Apache Iceberg
tables on the native Spark pipeline engine with ACID commit semantics.
-:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=yes, Beam
Spark=yes, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/spark-lake-table-output.svg[Spark lake table output
transform Icon, role="image-doc-icon"] Spark lake table output
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-sql.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-sql.adoc
index fc0f71514c..07114b9001 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-sql.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/spark-sql.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: Spark SQL runs a SQL statement over the Datasets of the incoming
transforms on the native Spark pipeline engine, so joins, unions and
aggregations are planned as one Catalyst query.
-:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=yes, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=no, Single Threaded=no, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/spark-sql.svg[Spark SQL transform Icon,
role="image-doc-icon"] Spark SQL
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerows.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerows.adoc
index 2498139e2a..bb22f56780 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerows.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerows.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: The Unique Rows transform removes duplicate rows from the
(sorted) input stream(s).
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=yes, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=native, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/uniquerows.svg[Unique Rows transform Icon,
role="image-doc-icon"] Unique Rows
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/transform/BaseTransform.java
b/engine/src/main/java/org/apache/hop/pipeline/transform/BaseTransform.java
index 235fa63029..fa729591b6 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/transform/BaseTransform.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/transform/BaseTransform.java
@@ -261,6 +261,9 @@ public class BaseTransform<Meta extends ITransformMeta,
Data extends ITransformD
/** The list of IRowListener interfaces */
protected List<IRowListener> rowListeners;
+ /** Read-only view of {@link #rowListeners}, created once: getRowListeners()
runs per row. */
+ private List<IRowListener> rowListenersView;
+
/** The list of destination-aware IRowToListener interfaces (target hops /
putRowTo) */
protected List<IRowToListener> rowToListeners;
@@ -440,6 +443,7 @@ public class BaseTransform<Meta extends ITransformMeta,
Data extends ITransformD
rowDistribution = transformMeta.getRowDistribution();
rowListeners = new CopyOnWriteArrayList<>();
+ rowListenersView = Collections.unmodifiableList(rowListeners);
rowToListeners = new CopyOnWriteArrayList<>();
resultFiles = new HashMap<>();
resultFilesLock = new ReentrantReadWriteLock();
@@ -3377,7 +3381,10 @@ public class BaseTransform<Meta extends ITransformMeta,
Data extends ITransformD
*/
@Override
public List<IRowListener> getRowListeners() {
- return Collections.unmodifiableList(rowListeners);
+ if (rowListenersView == null) {
+ rowListenersView = Collections.unmodifiableList(rowListeners);
+ }
+ return rowListenersView;
}
@Override
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
index f4c0377789..262cbaf111 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
@@ -20,10 +20,12 @@ package org.apache.hop.spark.core;
import java.io.File;
import java.io.Serializable;
import java.net.InetAddress;
+import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;
+import java.util.NoSuchElementException;
import org.apache.commons.lang3.StringUtils;
import org.apache.hop.core.Const;
import org.apache.hop.core.exception.HopException;
@@ -303,6 +305,8 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
? JsonRowMeta.fromJson(inputRowMetaJson)
: new org.apache.hop.core.row.RowMeta();
IRowMeta outputRowMeta = JsonRowMeta.fromJson(outputRowMetaJson);
+ final HopSparkRowConverter.RowCodec inputCodec =
+ HopSparkRowConverter.RowCodec.of(inputRowMeta);
PipelineMeta pipelineMeta = new PipelineMeta();
pipelineMeta.setName(transformName);
@@ -478,7 +482,17 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
List<Object[]> resultRows = new ArrayList<>();
List<List<Object[]>> targetResultRowsList = new ArrayList<>();
final boolean multiTarget = !targetTransforms.isEmpty();
- if (!multiTarget) {
+ // Plain case: one main input, no info or target streams — rows go
through a one-slot
+ // handler straight into the Spark output queue; no Injector, row sets
or executor loop.
+ final boolean direct =
+ !inputTransform
+ && !acceptFilenamesFromMain
+ && !multiTarget
+ && infoTransforms.isEmpty()
+ && mainTransform instanceof BaseTransform;
+ if (direct) {
+ // output captured by SparkDirectRowHandler.putRow
+ } else if (!multiTarget) {
transformCombi.transform.addRowListener(
new RowAdapter() {
@Override
@@ -533,27 +547,39 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
infoCombi.transform.processRow();
}
- List<Row> output = new ArrayList<>();
MetricsThrottle throttle = new MetricsThrottle();
final SparkTransformExecutionSampling samplingRef = sampling;
+ final ITransform mainRef = mainTransform;
+ final RowProducer mainProducer = rowProducer;
+ final TransformMetaDataCombi injectorCombi = mainInjectorCombi;
+ final long[] rowsSeen = {0L};
+ // Rows are handed to Spark lazily: every step() pushes at most one
input row (or runs one
+ // iteration of a source) and queues what the transform produced. The
partition output is
+ // never materialised, so memory is bounded by the transform's own
state, not by the
+ // partition size, and the downstream Spark operator can start consuming
immediately.
+ PartitionIterator rows;
if (inputTransform) {
// Source transform: drive until finished with no external input
- driveUntilDone(
- pipeline,
- executor,
- output,
- outputRowMeta,
- multiTarget,
- resultRows,
- targetTransforms,
- targetResultRowsList,
- throttle,
- mainTransform,
- copyNr,
- host,
- partitionStartMs,
- samplingRef);
+ rows =
+ new PartitionIterator(
+ outputRowMeta, multiTarget, targetTransforms, resultRows,
targetResultRowsList) {
+ private boolean more = true;
+
+ @Override
+ boolean step() throws HopException {
+ if (!more || pipeline.isFinished() || pipeline.getErrors() !=
0) {
+ return false;
+ }
+ clearCapture(resultRows, targetResultRowsList);
+ more = executor.oneIteration();
+ if (throttle.shouldPublish(resultRows.size())) {
+ publishMetrics(mainRef, copyNr, host, partitionStartMs,
true, false);
+ flushSamplesQuietly(samplingRef, false);
+ }
+ return true;
+ }
+ };
} else if (acceptFilenamesFromMain) {
// Accept-filenames (Text File Input, Excel, …):
// 1) Pre-load every partition row (filenames) into the main injector
and mark finished.
@@ -564,7 +590,6 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
// remaining hop size. After the first processRow consumes all
filenames, size is 0
// and further oneIteration() calls never advance the reader (infinite
stall after
// openNextFile/createReader — the log line users see just before
hang).
- long rowsSeen = 0;
while (input.hasNext()) {
Row sparkRow = input.next();
Object[] hopRow = HopSparkRowConverter.toHopRow(inputRowMeta,
sparkRow);
@@ -572,7 +597,7 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
if (mainInjectorCombi != null) {
mainInjectorCombi.transform.processRow();
}
- rowsSeen++;
+ rowsSeen[0]++;
}
if (rowProducer != null) {
rowProducer.finished();
@@ -580,82 +605,148 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
mainInjectorCombi.transform.processRow();
}
}
- if (rowsSeen > 0 || throttle.shouldPublish(0)) {
+ if (rowsSeen[0] > 0 || throttle.shouldPublish(0)) {
publishMetrics(mainTransform, copyNr, host, partitionStartMs, true,
false);
}
- driveAcceptFilenamesUntilDone(
- pipeline,
- mainTransform,
- output,
- outputRowMeta,
- multiTarget,
- resultRows,
- targetTransforms,
- targetResultRowsList,
- throttle,
- copyNr,
- host,
- partitionStartMs,
- samplingRef);
+ rows =
+ new PartitionIterator(
+ outputRowMeta, multiTarget, targetTransforms, resultRows,
targetResultRowsList) {
+ private boolean more = true;
+ private long contentRows = 0;
+
+ @Override
+ boolean step() throws HopException {
+ if (!more || pipeline.isFinished() || pipeline.getErrors() !=
0) {
+ return false;
+ }
+ clearCapture(resultRows, targetResultRowsList);
+ more = mainRef.processRow();
+ contentRows += resultRows.size();
+ if (throttle.shouldPublish(resultRows.size())
+ || contentRows % METRICS_ROW_INTERVAL == 0) {
+ publishMetrics(mainRef, copyNr, host, partitionStartMs,
true, false);
+ flushSamplesQuietly(samplingRef, false);
+ }
+ return true;
+ }
+ };
+ } else if (direct) {
+ final SparkDirectRowHandler[] handlerRef = new
SparkDirectRowHandler[1];
+ rows =
+ new PartitionIterator(
+ outputRowMeta, multiTarget, targetTransforms, resultRows,
targetResultRowsList) {
+ private boolean inputDrained = false;
+
+ @Override
+ boolean step() throws HopException {
+ SparkDirectRowHandler handler = handlerRef[0];
+ if (input.hasNext()) {
+ handler.offer(inputCodec.toHop(input.next()));
+ mainRef.processRow();
+ if (mainRef.getErrors() > 0) {
+ return false;
+ }
+ rowsSeen[0]++;
+ if (throttle.shouldPublish(1) || rowsSeen[0] %
METRICS_ROW_INTERVAL == 0) {
+ publishMetrics(mainRef, copyNr, host, partitionStartMs,
true, false);
+ flushSamplesQuietly(samplingRef, false);
+ }
+ return true;
+ }
+ if (inputDrained) {
+ return false;
+ }
+ // End of input: one call with an empty slot so getRow()
returns null and
+ // buffering transforms emit what they hold.
+ inputDrained = true;
+ handler.offer(null);
+ mainRef.processRow();
+ return true;
+ }
+ };
+ SparkDirectRowHandler handler =
+ new SparkDirectRowHandler(
+ (BaseTransform) mainTransform, inputRowMeta, outputRowMeta,
rows::emit);
+ ((BaseTransform) mainTransform).setRowHandler(handler);
+ handlerRef[0] = handler;
} else {
- long rowsSeen = 0;
- while (input.hasNext()) {
- Row sparkRow = input.next();
- Object[] hopRow = HopSparkRowConverter.toHopRow(inputRowMeta,
sparkRow);
- clearCapture(resultRows, targetResultRowsList);
- rowProducer.putRow(inputRowMeta, hopRow, false);
- // Forward the main row onto the hop before Stream Lookup's
info-first processRow
- // calls getRow() for the main stream (same timing Beam gets from
topo order +
- // non-blocking handler).
- if (mainInjectorCombi != null) {
- mainInjectorCombi.transform.processRow();
- }
- executor.oneIteration();
- appendCapturedRows(
- output,
- outputRowMeta,
- multiTarget,
- resultRows,
- targetTransforms,
- targetResultRowsList);
- rowsSeen++;
- if (throttle.shouldPublish(1) || rowsSeen % METRICS_ROW_INTERVAL ==
0) {
- publishMetrics(mainTransform, copyNr, host, partitionStartMs,
true, false);
- flushSamplesQuietly(samplingRef, false);
- }
- }
- if (rowProducer != null) {
- rowProducer.finished();
- if (mainInjectorCombi != null) {
- mainInjectorCombi.transform.processRow();
- }
- clearCapture(resultRows, targetResultRowsList);
- executor.oneIteration();
- appendCapturedRows(
- output,
- outputRowMeta,
- multiTarget,
- resultRows,
- targetTransforms,
- targetResultRowsList);
- }
- }
+ rows =
+ new PartitionIterator(
+ outputRowMeta, multiTarget, targetTransforms, resultRows,
targetResultRowsList) {
+ private boolean inputDrained = false;
- if (pipeline.getErrors() > 0) {
- publishMetrics(mainTransform, copyNr, host, partitionStartMs, false,
true);
- throw new HopException(
- "Errors detected while executing transform '"
- + transformName
- + "' on a Spark partition");
+ @Override
+ boolean step() throws HopException {
+ if (input.hasNext()) {
+ Object[] hopRow = inputCodec.toHop(input.next());
+ clearCapture(resultRows, targetResultRowsList);
+ mainProducer.putRow(inputRowMeta, hopRow, false);
+ // Forward the main row onto the hop before Stream Lookup's
info-first
+ // processRow calls getRow() for the main stream (same
timing Beam gets from
+ // topo order + non-blocking handler).
+ if (injectorCombi != null) {
+ injectorCombi.transform.processRow();
+ }
+ executor.oneIteration();
+ rowsSeen[0]++;
+ if (throttle.shouldPublish(1) || rowsSeen[0] %
METRICS_ROW_INTERVAL == 0) {
+ publishMetrics(mainRef, copyNr, host, partitionStartMs,
true, false);
+ flushSamplesQuietly(samplingRef, false);
+ }
+ return true;
+ }
+ if (inputDrained) {
+ return false;
+ }
+ // End of input: flag the injector done and give buffering
transforms one last
+ // iteration to emit what they hold.
+ inputDrained = true;
+ if (mainProducer != null) {
+ mainProducer.finished();
+ if (injectorCombi != null) {
+ injectorCombi.transform.processRow();
+ }
+ clearCapture(resultRows, targetResultRowsList);
+ executor.oneIteration();
+ }
+ return true;
+ }
+ };
}
- executor.dispose();
- publishMetrics(mainTransform, copyNr, host, partitionStartMs, false,
true);
- if (sampling != null) {
- sampling.close();
- sampling = null;
+ // End-of-partition duties run when the consumer exhausts the iterator,
or on task
+ // completion if it stops early (limit, cancelled stage).
+ final SparkTransformExecutionSampling samplingToClose = sampling;
+ Runnable finish =
+ () -> {
+ try {
+ if (pipeline.getErrors() > 0) {
+ publishMetrics(mainRef, copyNr, host, partitionStartMs, false,
true);
+ throw new HopException(
+ "Errors detected while executing transform '"
+ + transformName
+ + "' on a Spark partition");
+ }
+ executor.dispose();
+ publishMetrics(mainRef, copyNr, host, partitionStartMs, false,
true);
+ if (samplingToClose != null) {
+ samplingToClose.close();
+ }
+ } catch (HopException e) {
+ throw new HopRuntimeException(
+ "Error executing Hop transform '" + transformName + "' in
Spark mapPartitions",
+ e);
+ }
+ };
+ rows.onFinish(finish);
+ TaskContext taskContext = TaskContext.get();
+ if (taskContext != null) {
+ taskContext.addTaskCompletionListener(
+ (org.apache.spark.util.TaskCompletionListener) ctx ->
rows.finishQuietly());
}
- return output.iterator();
+ // Sampling is closed by finish(); the catch block below only handles
setup failures
+ sampling = null;
+ return rows;
} catch (Exception e) {
if (mainTransform != null) {
try {
@@ -996,99 +1087,112 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
* Drive the single-threaded mini-pipeline until it finishes (source
transforms with no main input
* hop — Row Generator, Get File Names, etc.).
*/
- private void driveUntilDone(
- LocalPipelineEngine pipeline,
- SingleThreadedPipelineExecutor executor,
- List<Row> output,
- IRowMeta outputRowMeta,
- boolean multiTarget,
- List<Object[]> resultRows,
- List<String> targetTransforms,
- List<List<Object[]>> targetResultRowsList,
- MetricsThrottle throttle,
- ITransform mainTransform,
- int copyNr,
- String host,
- long partitionStartMs,
- SparkTransformExecutionSampling samplingRef)
- throws HopException {
- boolean more = true;
- while (more && !pipeline.isFinished() && pipeline.getErrors() == 0) {
- clearCapture(resultRows, targetResultRowsList);
- more = executor.oneIteration();
- appendCapturedRows(
- output, outputRowMeta, multiTarget, resultRows, targetTransforms,
targetResultRowsList);
- if (throttle.shouldPublish(resultRows.size())) {
- publishMetrics(mainTransform, copyNr, host, partitionStartMs, true,
false);
- flushSamplesQuietly(samplingRef, false);
- }
+ private static void clearCapture(
+ List<Object[]> resultRows, List<List<Object[]>> targetResultRowsList) {
+ resultRows.clear();
+ for (List<Object[]> list : targetResultRowsList) {
+ list.clear();
}
}
/**
- * After filename rows are pre-loaded onto the main injector, call {@code
processRow()} on the
- * file-input transform until it finishes reading content.
- *
- * <p>{@link SingleThreadedPipelineExecutor#oneIteration()} only schedules
{@code processRow} once
- * per remaining input-rowset size when the transform has input hops. After
accept-filenames
- * drains those rows, size is 0 and further iterations never advance the
reader.
+ * Lazy partition output. {@link #step()} advances the mini-pipeline by one
unit of work and
+ * leaves whatever the transform emitted in the capture lists; those rows
are converted and
+ * queued, and handed to Spark one at a time.
*/
- private void driveAcceptFilenamesUntilDone(
- LocalPipelineEngine pipeline,
- ITransform mainTransform,
- List<Row> output,
- IRowMeta outputRowMeta,
- boolean multiTarget,
- List<Object[]> resultRows,
- List<String> targetTransforms,
- List<List<Object[]>> targetResultRowsList,
- MetricsThrottle throttle,
- int copyNr,
- String host,
- long partitionStartMs,
- SparkTransformExecutionSampling samplingRef)
- throws HopException {
- boolean more = true;
- long contentRows = 0;
- while (more && !pipeline.isFinished() && pipeline.getErrors() == 0) {
- clearCapture(resultRows, targetResultRowsList);
- more = mainTransform.processRow();
- appendCapturedRows(
- output, outputRowMeta, multiTarget, resultRows, targetTransforms,
targetResultRowsList);
- contentRows += resultRows.size();
- if (throttle.shouldPublish(resultRows.size()) || contentRows %
METRICS_ROW_INTERVAL == 0) {
- publishMetrics(mainTransform, copyNr, host, partitionStartMs, true,
false);
- flushSamplesQuietly(samplingRef, false);
+ private abstract static class PartitionIterator implements Iterator<Row> {
+ private final HopSparkRowConverter.RowCodec outputCodec;
+ private final boolean multiTarget;
+ private final List<String> targetTransforms;
+ private final List<Object[]> resultRows;
+ private final List<List<Object[]>> targetResultRowsList;
+ private final ArrayDeque<Row> pending = new ArrayDeque<>();
+ private Runnable finish;
+ private boolean finished;
+ private boolean exhausted;
+
+ PartitionIterator(
+ IRowMeta outputRowMeta,
+ boolean multiTarget,
+ List<String> targetTransforms,
+ List<Object[]> resultRows,
+ List<List<Object[]>> targetResultRowsList) {
+ this.outputCodec = HopSparkRowConverter.RowCodec.of(outputRowMeta);
+ this.multiTarget = multiTarget;
+ this.targetTransforms = targetTransforms;
+ this.resultRows = resultRows;
+ this.targetResultRowsList = targetResultRowsList;
+ }
+
+ /** Do one unit of work; return false when the transform has nothing left
to produce. */
+ abstract boolean step() throws HopException;
+
+ /** Queue an already converted output row (direct handler path). */
+ void emit(Row row) {
+ pending.add(row);
+ }
+
+ void onFinish(Runnable finish) {
+ this.finish = finish;
+ }
+
+ @Override
+ public boolean hasNext() {
+ try {
+ while (pending.isEmpty() && !exhausted) {
+ if (step()) {
+ queueCaptured();
+ } else {
+ exhausted = true;
+ }
+ }
+ } catch (Exception e) {
+ throw new HopRuntimeException("Error executing Hop transform in Spark
mapPartitions", e);
+ }
+ if (pending.isEmpty()) {
+ runFinish();
+ return false;
}
+ return true;
}
- }
- private static void clearCapture(
- List<Object[]> resultRows, List<List<Object[]>> targetResultRowsList) {
- resultRows.clear();
- for (List<Object[]> list : targetResultRowsList) {
- list.clear();
+ @Override
+ public Row next() {
+ if (!hasNext()) {
+ throw new NoSuchElementException();
+ }
+ return pending.poll();
}
- }
- private static void appendCapturedRows(
- List<Row> output,
- IRowMeta outputRowMeta,
- boolean multiTarget,
- List<Object[]> resultRows,
- List<String> targetTransforms,
- List<List<Object[]>> targetResultRowsList)
- throws HopException {
- if (!multiTarget) {
- for (Object[] hopRow : resultRows) {
- output.add(HopSparkRowConverter.toSparkRow(outputRowMeta, hopRow));
+ private void queueCaptured() throws HopException {
+ if (!multiTarget) {
+ for (Object[] hopRow : resultRows) {
+ pending.add(outputCodec.toSpark(hopRow));
+ }
+ return;
+ }
+ for (int t = 0; t < targetResultRowsList.size(); t++) {
+ for (Object[] hopRow : targetResultRowsList.get(t)) {
+ pending.add(outputCodec.toTaggedSpark(targetTransforms.get(t),
hopRow));
+ }
+ }
+ }
+
+ private void runFinish() {
+ if (!finished) {
+ finished = true;
+ if (finish != null) {
+ finish.run();
+ }
}
- return;
}
- for (int t = 0; t < targetTransforms.size(); t++) {
- String tag = targetTransforms.get(t);
- for (Object[] hopRow : targetResultRowsList.get(t)) {
- output.add(HopSparkRowConverter.toTaggedSparkRow(tag, outputRowMeta,
hopRow));
+
+ /** Task-completion hook: release the mini-pipeline even when Spark
stopped pulling early. */
+ void finishQuietly() {
+ try {
+ runFinish();
+ } catch (Exception e) {
+ LogChannel.GENERAL.logError("Error finishing Hop transform partition
(non-fatal)", e);
}
}
}
@@ -1098,7 +1202,14 @@ public class HopMapPartitionsFn implements
MapPartitionsFunction<Row, Row>, Seri
private static final long serialVersionUID = 1L;
private long lastPublishMs = 0L;
+ private int calls = 0;
+
boolean shouldPublish(int rowsThisBatch) {
+ // The clock read is a syscall; sampling it every 128 calls keeps the ~1
s cadence for slow
+ // transforms without paying for it on every row of a fast one.
+ if (lastPublishMs != 0L && (++calls & 127) != 0) {
+ return false;
+ }
long now = System.currentTimeMillis();
if (lastPublishMs == 0L || now - lastPublishMs >=
METRICS_TIME_INTERVAL_MS) {
lastPublishMs = now;
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopSparkRowConverter.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopSparkRowConverter.java
index 55a0a462b4..e1d404c848 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopSparkRowConverter.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopSparkRowConverter.java
@@ -77,6 +77,97 @@ public final class HopSparkRowConverter {
};
}
+ /**
+ * Per-partition converter: resolves the value metas and type ids once, then
converts rows with a
+ * flat switch and direct pass-through for values that already have the
Spark-side Java type
+ * (String, Long, Double, Boolean). Falls back to the generic per-value
conversion for the rest.
+ */
+ public static final class RowCodec {
+ private final IValueMeta[] metas;
+ private final int[] types;
+ private final boolean[] plain;
+
+ private RowCodec(IRowMeta rowMeta) {
+ int n = rowMeta.size();
+ metas = new IValueMeta[n];
+ types = new int[n];
+ plain = new boolean[n];
+ for (int i = 0; i < n; i++) {
+ metas[i] = rowMeta.getValueMeta(i);
+ types[i] = metas[i].getType();
+ plain[i] = metas[i].getStorageType() == IValueMeta.STORAGE_TYPE_NORMAL;
+ }
+ }
+
+ public static RowCodec of(IRowMeta rowMeta) {
+ return new RowCodec(rowMeta);
+ }
+
+ public int size() {
+ return metas.length;
+ }
+
+ public Object[] toHop(Row sparkRow) throws HopException {
+ Object[] hopRow = new Object[metas.length];
+ if (sparkRow == null) {
+ return hopRow;
+ }
+ int available = Math.min(metas.length, sparkRow.length());
+ for (int i = 0; i < available; i++) {
+ Object v = sparkRow.get(i);
+ if (v == null) {
+ continue;
+ }
+ hopRow[i] =
+ switch (types[i]) {
+ case IValueMeta.TYPE_STRING -> v instanceof String ? v :
v.toString();
+ case IValueMeta.TYPE_INTEGER -> v instanceof Long ? v :
toHopValue(metas[i], v);
+ case IValueMeta.TYPE_NUMBER -> v instanceof Double ? v :
toHopValue(metas[i], v);
+ case IValueMeta.TYPE_BOOLEAN -> v instanceof Boolean ? v :
toHopValue(metas[i], v);
+ default -> toHopValue(metas[i], v);
+ };
+ }
+ return hopRow;
+ }
+
+ public Row toSpark(Object[] hopRow) throws HopException {
+ return RowFactory.create(sparkValues(hopRow, 0));
+ }
+
+ public Row toTaggedSpark(String targetTag, Object[] hopRow) throws
HopException {
+ Object[] values = sparkValues(hopRow, 1);
+ values[0] = targetTag != null ? targetTag : HopSparkUtil.MAIN_TARGET_TAG;
+ return RowFactory.create(values);
+ }
+
+ private Object[] sparkValues(Object[] hopRow, int offset) throws
HopException {
+ Object[] values = new Object[metas.length + offset];
+ int available = hopRow == null ? 0 : Math.min(metas.length,
hopRow.length);
+ for (int i = 0; i < available; i++) {
+ Object v = hopRow[i];
+ if (v == null) {
+ continue;
+ }
+ Object out;
+ if (plain[i]) {
+ out =
+ switch (types[i]) {
+ case IValueMeta.TYPE_STRING -> v instanceof String ? v :
toSparkValue(metas[i], v);
+ case IValueMeta.TYPE_INTEGER -> v instanceof Long ? v :
toSparkValue(metas[i], v);
+ case IValueMeta.TYPE_NUMBER -> v instanceof Double ? v :
toSparkValue(metas[i], v);
+ case IValueMeta.TYPE_BOOLEAN ->
+ v instanceof Boolean ? v : toSparkValue(metas[i], v);
+ default -> toSparkValue(metas[i], v);
+ };
+ } else {
+ out = toSparkValue(metas[i], v);
+ }
+ values[i + offset] = out;
+ }
+ return values;
+ }
+ }
+
public static Row toSparkRow(IRowMeta rowMeta, Object[] hopRow) throws
HopException {
Object[] values = new Object[rowMeta.size()];
for (int i = 0; i < rowMeta.size(); i++) {
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkDirectRowHandler.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkDirectRowHandler.java
new file mode 100644
index 0000000000..44e8ceea1c
--- /dev/null
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkDirectRowHandler.java
@@ -0,0 +1,101 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.spark.core;
+
+import java.util.List;
+import java.util.function.Consumer;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.exception.HopTransformException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.pipeline.transform.BaseTransform;
+import org.apache.hop.pipeline.transform.IRowListener;
+import org.apache.spark.sql.Row;
+
+/**
+ * Row handler for the plain generic case (one main input, no info or target
streams): the partition
+ * loop hands the current input row to the transform through a one-slot buffer
and takes what the
+ * transform writes straight into the Spark output queue, already converted.
+ *
+ * <p>This keeps the per-row protocol of the row-set based path — {@code
processRow()} is called
+ * once per input row and once more with an empty slot at the end so {@code
getRow()} returns null —
+ * but skips the Injector transform, the two row sets and the executor
bookkeeping in between.
+ */
+public class SparkDirectRowHandler extends SparkRowHandler {
+
+ private final BaseTransform transform;
+ private final IRowMeta inputRowMeta;
+ private final IRowMeta outputRowMeta;
+ private final HopSparkRowConverter.RowCodec outputCodec;
+ private final Consumer<Row> output;
+ private Object[] slot;
+ private boolean inputRowMetaSet;
+
+ public SparkDirectRowHandler(
+ BaseTransform transform,
+ IRowMeta inputRowMeta,
+ IRowMeta outputRowMeta,
+ Consumer<Row> output) {
+ super(transform);
+ this.transform = transform;
+ this.inputRowMeta = inputRowMeta;
+ this.outputRowMeta = outputRowMeta;
+ this.outputCodec = HopSparkRowConverter.RowCodec.of(outputRowMeta);
+ this.output = output;
+ }
+
+ /** Make {@code row} the next row {@link #getRow()} returns; null means end
of input. */
+ public void offer(Object[] row) {
+ this.slot = row;
+ }
+
+ @Override
+ public Object[] getRow() throws HopException {
+ Object[] row = slot;
+ slot = null;
+ if (!inputRowMetaSet) {
+ transform.setInputRowMeta(inputRowMeta);
+ inputRowMetaSet = true;
+ }
+ if (row != null) {
+ transform.incrementLinesRead();
+ List<IRowListener> rowListeners = transform.getRowListeners();
+ if (!rowListeners.isEmpty()) {
+ for (IRowListener rowListener : rowListeners) {
+ rowListener.rowReadEvent(inputRowMeta, row);
+ }
+ }
+ }
+ return row;
+ }
+
+ @Override
+ public void putRow(IRowMeta rowMeta, Object[] row) throws
HopTransformException {
+ List<IRowListener> rowListeners = transform.getRowListeners();
+ if (!rowListeners.isEmpty()) {
+ for (IRowListener rowListener : rowListeners) {
+ rowListener.rowWrittenEvent(rowMeta, row);
+ }
+ }
+ try {
+ output.accept(outputCodec.toSpark(row));
+ } catch (HopException e) {
+ throw new HopTransformException(e);
+ }
+ transform.incrementLinesWritten();
+ }
+}
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetrics.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetrics.java
index 0c19f06a8f..c33cbff719 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetrics.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetrics.java
@@ -17,23 +17,21 @@
package org.apache.hop.spark.core;
-import java.io.Serializable;
-import java.net.InetAddress;
-import java.util.Iterator;
-import java.util.Objects;
-import org.apache.spark.TaskContext;
-import org.apache.spark.api.java.JavaRDD;
+import static org.apache.spark.sql.functions.count;
+import static org.apache.spark.sql.functions.lit;
+
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
-import org.apache.spark.sql.types.StructType;
/**
* Instruments native Spark {@link Dataset} stages so row flow is visible in
Hop {@code
- * EngineMetrics}. Inserts a pass-through {@code mapPartitions} that reports
absolute counters into
- * the same {@link SparkTransformMetricsAccumulator} used by {@link
HopMapPartitionsFn}.
+ * EngineMetrics}.
*
- * <p>Counters are only updated when the lineage is materialized (an action).
Per-partition {@code
- * copyNr} matches {@link TaskContext#partitionId()}.
+ * <p>Each tracked stage gets a named {@code CollectMetrics} node ({@link
Dataset#observe}) that
+ * stays inside the Catalyst plan: no RDD round-trip, so whole-stage codegen,
column pruning and
+ * filter pushdown are preserved. The counts themselves are read from Spark's
own per-task SQL
+ * metrics by {@link SparkNativeMetricsListener}, which attributes them to the
transform via the
+ * anchor and reports partition / host / timing from the task info.
*/
public final class SparkNativeMetrics {
@@ -47,164 +45,18 @@ public final class SparkNativeMetrics {
TRANSFORM
}
- private static final int ROW_INTERVAL = 1000;
- private static final long TIME_INTERVAL_MS = 1000L;
-
private SparkNativeMetrics() {}
/**
- * Wrap {@code dataset} so that when it is computed, each partition reports
metrics for {@code
- * transformName}. Returns {@code dataset} unchanged when accumulator or
name is null.
+ * Anchor {@code dataset} so that its row count is attributed to {@code
transformName} by the
+ * listener. Returns {@code dataset} unchanged when listener or name is null.
*/
public static Dataset<Row> track(
- Dataset<Row> dataset,
- String transformName,
- SparkTransformMetricsAccumulator accumulator,
- Role role) {
- if (dataset == null
- || accumulator == null
- || transformName == null
- || transformName.isEmpty()) {
+ Dataset<Row> dataset, String transformName, SparkNativeMetricsListener
listener, Role role) {
+ if (dataset == null || listener == null || transformName == null ||
transformName.isEmpty()) {
return dataset;
}
- Role effectiveRole = role != null ? role : Role.TRANSFORM;
- StructType schema = dataset.schema();
- JavaRDD<Row> tracked =
- dataset
- .toJavaRDD()
- .mapPartitions(
- iterator ->
- new CountingIterator(iterator, transformName, accumulator,
effectiveRole),
- true);
- return dataset.sparkSession().createDataFrame(tracked, schema);
- }
-
- static final class CountingIterator implements Iterator<Row>, Serializable {
- private static final long serialVersionUID = 1L;
-
- private final Iterator<Row> source;
- private final String transformName;
- private final SparkTransformMetricsAccumulator accumulator;
- private final Role role;
- private final int copyNr;
- private final String host;
-
- private long count;
- private long partitionStartMs;
- private long lastPublishMs;
- private boolean finishedPublished;
- private boolean startedPublished;
-
- CountingIterator(
- Iterator<Row> source,
- String transformName,
- SparkTransformMetricsAccumulator accumulator,
- Role role) {
- this.source = Objects.requireNonNull(source, "source");
- this.transformName = transformName;
- this.accumulator = accumulator;
- this.role = role;
- this.copyNr = partitionId();
- this.host = localHost();
- }
-
- @Override
- public boolean hasNext() {
- ensureStarted();
- if (source.hasNext()) {
- return true;
- }
- publishFinished();
- return false;
- }
-
- @Override
- public Row next() {
- ensureStarted();
- Row row = source.next();
- count++;
- if (shouldPublishProgress()) {
- publish(true, false);
- }
- return row;
- }
-
- private void ensureStarted() {
- if (!startedPublished) {
- startedPublished = true;
- partitionStartMs = System.currentTimeMillis();
- lastPublishMs = partitionStartMs;
- publish(true, false);
- }
- }
-
- private boolean shouldPublishProgress() {
- if (count > 0 && count % ROW_INTERVAL == 0) {
- return true;
- }
- long now = System.currentTimeMillis();
- if (now - lastPublishMs >= TIME_INTERVAL_MS) {
- lastPublishMs = now;
- return true;
- }
- return false;
- }
-
- private void publishFinished() {
- if (!finishedPublished) {
- finishedPublished = true;
- publish(false, true);
- }
- }
-
- private void publish(boolean running, boolean finished) {
- long read = 0;
- long written = 0;
- long input = 0;
- long output = 0;
- switch (role) {
- case INPUT:
- input = count;
- written = count;
- break;
- case OUTPUT:
- read = count;
- output = count;
- break;
- case TRANSFORM:
- default:
- read = count;
- written = count;
- break;
- }
- long endTimeMs = finished ? System.currentTimeMillis() : 0L;
- accumulator.add(
- new SparkTransformMetricSlice(
- transformName,
- copyNr,
- host,
- read,
- written,
- input,
- output,
- 0,
- running,
- finished,
- partitionStartMs,
- endTimeMs));
- }
-
- private static int partitionId() {
- TaskContext ctx = TaskContext.get();
- return ctx != null ? ctx.partitionId() : 0;
- }
-
- private static String localHost() {
- try {
- return InetAddress.getLocalHost().getHostName();
- } catch (Exception e) {
- return null;
- }
- }
+ String token = listener.register(transformName, role != null ? role :
Role.TRANSFORM);
+ return dataset.observe(token, count(lit(1)).alias("rows"));
}
}
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetricsListener.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetricsListener.java
new file mode 100644
index 0000000000..72eac0d0d2
--- /dev/null
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/SparkNativeMetricsListener.java
@@ -0,0 +1,350 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.spark.core;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.spark.SparkContext;
+import org.apache.spark.scheduler.AccumulableInfo;
+import org.apache.spark.scheduler.SparkListener;
+import org.apache.spark.scheduler.SparkListenerEvent;
+import org.apache.spark.scheduler.SparkListenerExecutorMetricsUpdate;
+import org.apache.spark.scheduler.SparkListenerTaskEnd;
+import org.apache.spark.scheduler.SparkListenerTaskStart;
+import org.apache.spark.scheduler.TaskInfo;
+import org.apache.spark.sql.execution.SparkPlanInfo;
+import org.apache.spark.sql.execution.metric.SQLMetricInfo;
+import
org.apache.spark.sql.execution.ui.SparkListenerSQLAdaptiveExecutionUpdate;
+import org.apache.spark.sql.execution.ui.SparkListenerSQLExecutionStart;
+import scala.Option;
+import scala.Tuple4;
+import scala.jdk.javaapi.CollectionConverters;
+
+/**
+ * Driver-side {@link SparkListener} that turns Spark's own per-task SQL
metrics into Hop {@link
+ * SparkTransformMetricSlice}s for native Dataset stages.
+ *
+ * <p>{@link SparkNativeMetrics#track} inserts a named {@code CollectMetrics}
node ({@code
+ * Dataset.observe}) as an anchor for each transform. This listener resolves
that anchor in the
+ * physical plan to the accumulator that counts its rows, then reads the
per-task value of that
+ * accumulator from task events: the partition index, host and launch/finish
time come from the
+ * {@link TaskInfo}. Nothing is inserted into the Dataset lineage, so
whole-stage codegen and column
+ * pruning stay intact.
+ *
+ * <p>Row counts for nodes directly below the anchor come from the node's
{@code number of output
+ * rows} metric. Nodes that consume a shuffle (Sort, Coalesce, …) carry no
such metric, so their
+ * anchor is bound to the consuming stage instead (identified by the node's
other metric ids) and
+ * the task's shuffle-read record count is used. Slices are absolute snapshots
merged with max
+ * semantics in the accumulator, so a stage that runs twice (range-partition
sampling before a sort)
+ * does not double count.
+ */
+public class SparkNativeMetricsListener extends SparkListener {
+
+ static final String TOKEN_PREFIX = "hop_metrics_";
+ private static final String NUM_OUTPUT_ROWS = "number of output rows";
+ private static final String SHUFFLE_READ_RECORDS =
"internal.metrics.shuffle.read.recordsRead";
+ private static final Set<String> SHUFFLE_NODES =
+ Set.of(
+ "Exchange",
+ "ShuffleQueryStage",
+ "AQEShuffleRead",
+ "CustomShuffleReader",
+ "BroadcastExchange",
+ "BroadcastQueryStage");
+ private static final Set<String> PASS_THROUGH_NODES =
+ Set.of("WholeStageCodegen", "InputAdapter", "CollectMetrics",
"AdaptiveSparkPlan");
+
+ private final SparkTransformMetricsAccumulator sink;
+ private final String prefix;
+ private final AtomicInteger sequence = new AtomicInteger();
+
+ /** observe() token → registration. */
+ private final Map<String, Registration> registrations = new
ConcurrentHashMap<>();
+
+ /** accumulator id of a "number of output rows" metric → registrations
counted by it. */
+ private final Map<Long, List<Registration>> directIds = new
ConcurrentHashMap<>();
+
+ /** accumulator id of any metric of a shuffle-consuming node → registrations
counted by it. */
+ private final Map<Long, List<Registration>> shuffleMarkerIds = new
ConcurrentHashMap<>();
+
+ /** taskId → task info (for live updates that carry no TaskInfo). */
+ private final Map<Long, TaskInfo> tasks = new ConcurrentHashMap<>();
+
+ public SparkNativeMetricsListener(SparkTransformMetricsAccumulator sink) {
+ this.sink = sink;
+ this.prefix = TOKEN_PREFIX + UUID.randomUUID().toString().substring(0, 8)
+ "_";
+ }
+
+ /** Register a transform and return the unique observation name to anchor it
with. */
+ public String register(String transformName, SparkNativeMetrics.Role role) {
+ String token = prefix + sequence.incrementAndGet();
+ registrations.put(token, new Registration(transformName, role));
+ return token;
+ }
+
+ public void addTo(SparkContext sparkContext) {
+ sparkContext.addSparkListener(this);
+ }
+
+ public void removeFrom(SparkContext sparkContext) {
+ try {
+ sparkContext.removeSparkListener(this);
+ } catch (Exception ignored) {
+ // context may already be stopped
+ }
+ }
+
+ /** Block until queued listener events are delivered so final counts are
visible. */
+ public void flush(SparkContext sparkContext, long timeoutMs) {
+ try {
+ sparkContext.listenerBus().waitUntilEmpty(timeoutMs);
+ } catch (Exception ignored) {
+ // best effort; a late event only delays the last snapshot
+ }
+ }
+
+ // ---- plan resolution
--------------------------------------------------------------------
+
+ @Override
+ public void onOtherEvent(SparkListenerEvent event) {
+ if (event instanceof SparkListenerSQLExecutionStart start) {
+ resolvePlan(start.sparkPlanInfo());
+ } else if (event instanceof SparkListenerSQLAdaptiveExecutionUpdate
update) {
+ resolvePlan(update.sparkPlanInfo());
+ }
+ }
+
+ void resolvePlan(SparkPlanInfo root) {
+ if (root == null) {
+ return;
+ }
+ List<SparkPlanInfo> stack = new ArrayList<>();
+ stack.add(root);
+ while (!stack.isEmpty()) {
+ SparkPlanInfo node = stack.remove(stack.size() - 1);
+ Registration registration = anchorOf(node);
+ if (registration != null) {
+ bind(registration, node);
+ }
+ stack.addAll(children(node));
+ }
+ }
+
+ private Registration anchorOf(SparkPlanInfo node) {
+ if (!"CollectMetrics".equals(node.nodeName())) {
+ return null;
+ }
+ // simpleString: "CollectMetrics <name>, [<aggregates>]"
+ String s = node.simpleString();
+ int from = s.indexOf(prefix);
+ if (from < 0) {
+ return null;
+ }
+ int to = s.indexOf(',', from);
+ String token = to < 0 ? s.substring(from) : s.substring(from, to);
+ return registrations.get(token.trim());
+ }
+
+ /**
+ * Walk down the single-child chain below an anchor: bind to the first
"number of output rows"
+ * metric, or, when a shuffle boundary is reached first, to the metrics of
the consuming node.
+ */
+ private void bind(Registration registration, SparkPlanInfo anchor) {
+ List<Long> markers = new ArrayList<>();
+ SparkPlanInfo node = firstChild(anchor);
+ while (node != null) {
+ Long numOutputRows = metricId(node, NUM_OUTPUT_ROWS);
+ if (numOutputRows != null) {
+ directIds.computeIfAbsent(numOutputRows, k -> new
ArrayList<>()).add(registration);
+ return;
+ }
+ if (SHUFFLE_NODES.contains(node.nodeName())) {
+ if (markers.isEmpty()) {
+ // No consumer metrics between anchor and shuffle: attribute to the
map side instead
+ node = firstChild(node);
+ continue;
+ }
+ for (Long id : markers) {
+ shuffleMarkerIds.computeIfAbsent(id, k -> new
ArrayList<>()).add(registration);
+ }
+ return;
+ }
+ if (!isPassThrough(node)) {
+ for (SQLMetricInfo metric : metrics(node)) {
+ markers.add(metric.accumulatorId());
+ }
+ }
+ node = firstChild(node);
+ }
+ }
+
+ /** Wrapper nodes whose metrics (e.g. codegen "duration") say nothing about
row flow. */
+ private static boolean isPassThrough(SparkPlanInfo node) {
+ // nodeName carries a suffix for codegen stages: "WholeStageCodegen (2)"
+ for (String name : PASS_THROUGH_NODES) {
+ if (node.nodeName().startsWith(name)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static SparkPlanInfo firstChild(SparkPlanInfo node) {
+ List<SparkPlanInfo> children = children(node);
+ return children.size() == 1 ? children.get(0) : null;
+ }
+
+ private static List<SparkPlanInfo> children(SparkPlanInfo node) {
+ return CollectionConverters.asJava(node.children());
+ }
+
+ private static List<SQLMetricInfo> metrics(SparkPlanInfo node) {
+ return CollectionConverters.asJava(node.metrics());
+ }
+
+ private static Long metricId(SparkPlanInfo node, String name) {
+ for (SQLMetricInfo metric : metrics(node)) {
+ if (name.equals(metric.name())) {
+ return metric.accumulatorId();
+ }
+ }
+ return null;
+ }
+
+ // ---- task events
------------------------------------------------------------------------
+
+ @Override
+ public void onTaskStart(SparkListenerTaskStart taskStart) {
+ TaskInfo info = taskStart.taskInfo();
+ if (info != null) {
+ tasks.put(info.taskId(), info);
+ }
+ }
+
+ @Override
+ public void onExecutorMetricsUpdate(SparkListenerExecutorMetricsUpdate
update) {
+ for (Tuple4<Object, Object, Object,
scala.collection.immutable.Seq<AccumulableInfo>> entry :
+ CollectionConverters.asJava(update.accumUpdates())) {
+ long taskId = ((Number) entry._1()).longValue();
+ TaskInfo info = tasks.get(taskId);
+ if (info != null) {
+ report(info, CollectionConverters.asJava(entry._4()), false);
+ }
+ }
+ }
+
+ @Override
+ public void onTaskEnd(SparkListenerTaskEnd taskEnd) {
+ TaskInfo info = taskEnd.taskInfo();
+ if (info == null) {
+ return;
+ }
+ tasks.remove(info.taskId());
+ if (info.successful()) {
+ report(info, CollectionConverters.asJava(info.accumulables()), true);
+ }
+ }
+
+ private void report(TaskInfo info, List<AccumulableInfo> accumulables,
boolean finished) {
+ if (accumulables == null || accumulables.isEmpty()) {
+ return;
+ }
+ Map<Long, Long> byId = new HashMap<>();
+ long shuffleReadRecords = -1;
+ for (AccumulableInfo acc : accumulables) {
+ Long value = longValue(acc.update());
+ if (value == null) {
+ continue;
+ }
+ byId.put(acc.id(), value);
+ if (acc.name().isDefined() &&
SHUFFLE_READ_RECORDS.equals(acc.name().get())) {
+ shuffleReadRecords = value;
+ }
+ }
+ Map<Registration, Long> counts = new HashMap<>();
+ for (Map.Entry<Long, Long> e : byId.entrySet()) {
+ List<Registration> direct = directIds.get(e.getKey());
+ if (direct != null) {
+ for (Registration r : direct) {
+ counts.merge(r, e.getValue(), Math::max);
+ }
+ }
+ List<Registration> viaShuffle = shuffleMarkerIds.get(e.getKey());
+ if (viaShuffle != null && shuffleReadRecords >= 0) {
+ for (Registration r : viaShuffle) {
+ counts.merge(r, shuffleReadRecords, Math::max);
+ }
+ }
+ }
+ for (Map.Entry<Registration, Long> e : counts.entrySet()) {
+ sink.add(slice(e.getKey(), info, e.getValue(), finished));
+ }
+ }
+
+ private static Long longValue(Option<Object> update) {
+ if (update == null || update.isEmpty()) {
+ return null;
+ }
+ Object v = update.get();
+ return v instanceof Number n ? n.longValue() : null;
+ }
+
+ private static SparkTransformMetricSlice slice(
+ Registration registration, TaskInfo info, long count, boolean finished) {
+ long read = 0;
+ long written = 0;
+ long input = 0;
+ long output = 0;
+ switch (registration.role()) {
+ case INPUT -> {
+ input = count;
+ written = count;
+ }
+ case OUTPUT -> {
+ read = count;
+ output = count;
+ }
+ default -> {
+ read = count;
+ written = count;
+ }
+ }
+ return new SparkTransformMetricSlice(
+ registration.transformName(),
+ info.index(),
+ info.host(),
+ read,
+ written,
+ input,
+ output,
+ 0,
+ !finished,
+ finished,
+ info.launchTime(),
+ finished ? info.finishTime() : 0L);
+ }
+
+ /** A transform anchored by one observe() token. */
+ record Registration(String transformName, SparkNativeMetrics.Role role) {}
+}
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
index 8c32a8024c..a410b89c62 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
@@ -84,6 +84,7 @@ import
org.apache.hop.pipeline.engine.PipelineEngineCapabilities;
import org.apache.hop.pipeline.engine.PipelineEnginePlugin;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.core.SparkExecutionDataAccumulator;
+import org.apache.hop.spark.core.SparkNativeMetricsListener;
import org.apache.hop.spark.core.SparkTransformMetricSlice;
import org.apache.hop.spark.core.SparkTransformMetricsAccumulator;
import org.apache.hop.spark.execution.SparkTransformExecutionSampling;
@@ -173,6 +174,7 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
private Dataset<Row> resultDataset;
private Thread sparkThread;
private SparkTransformMetricsAccumulator metricsAccumulator;
+ private SparkNativeMetricsListener metricsListener;
private SparkExecutionDataAccumulator sampleDataAccumulator;
private Timer metricsRefreshTimer;
@@ -296,6 +298,10 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
metricsAccumulator = new SparkTransformMetricsAccumulator();
sparkSession.sparkContext().register(metricsAccumulator,
"hop-transform-metrics");
converter.setMetricsAccumulator(metricsAccumulator);
+ // Native Dataset stages: Spark's own SQL metrics → same accumulator,
via a listener
+ metricsListener = new SparkNativeMetricsListener(metricsAccumulator);
+ metricsListener.addTo(sparkSession.sparkContext());
+ converter.setMetricsListener(metricsListener);
sampleDataAccumulator = new SparkExecutionDataAccumulator();
sparkSession.sparkContext().register(sampleDataAccumulator,
"hop-execution-sample-data");
converter.setSampleDataAccumulator(sampleDataAccumulator);
@@ -350,6 +356,10 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
}
ExecutorUtil.cleanup(metricsRefreshTimer);
try {
+ if (metricsListener != null && sparkSession != null) {
+ // Drain the listener bus so the last task events are in
the accumulator
+ metricsListener.flush(sparkSession.sparkContext(),
5000L);
+ }
populateEngineMetrics();
} catch (Exception emEx) {
logChannel.logError("Error populating final engine
metrics", emEx);
@@ -1207,6 +1217,11 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
public void cleanup() {
ExecutorUtil.cleanup(metricsRefreshTimer);
ExecutorUtil.cleanup(executionInfoTimer);
+ if (sparkSession != null && metricsListener != null) {
+ // Shared / nested sessions keep running: never leave a stale listener
behind
+ metricsListener.removeFrom(sparkSession.sparkContext());
+ metricsListener = null;
+ }
if (sparkSession != null) {
try {
// Only stop sessions this engine created. Nested Pipeline Executor
children and
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
index 76298a394c..60bcf454dc 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
@@ -39,6 +39,7 @@ import org.apache.hop.pipeline.transform.BaseTransform;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.core.HopSparkUtil;
import org.apache.hop.spark.core.SparkExecutionDataAccumulator;
+import org.apache.hop.spark.core.SparkNativeMetricsListener;
import org.apache.hop.spark.core.SparkTransformMetricsAccumulator;
import org.apache.hop.spark.engines.ISparkPipelineEngineRunConfiguration;
import org.apache.hop.spark.pipeline.handler.SparkBaseTransformHandler;
@@ -492,6 +493,15 @@ public class HopPipelineMetaToSparkConverter {
}
}
+ /** Driver-side listener that attributes Spark's own SQL metrics to native
handler stages. */
+ public void setMetricsListener(SparkNativeMetricsListener metricsListener) {
+ for (ISparkPipelineTransformHandler handler : transformHandlers.values()) {
+ if (handler instanceof SparkBaseTransformHandler baseHandler) {
+ baseHandler.setMetricsListener(metricsListener);
+ }
+ }
+ }
+
public SparkTransformMetricsAccumulator getMetricsAccumulator() {
return metricsAccumulator;
}
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkBaseTransformHandler.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkBaseTransformHandler.java
index 2ee63a591e..c0e4ad1cb3 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkBaseTransformHandler.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkBaseTransformHandler.java
@@ -24,6 +24,7 @@ import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.ITransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.spark.core.SparkNativeMetrics;
+import org.apache.hop.spark.core.SparkNativeMetricsListener;
import org.apache.hop.spark.core.SparkTransformMetricsAccumulator;
import org.apache.hop.spark.pipeline.ISparkPipelineTransformHandler;
import org.apache.spark.sql.Dataset;
@@ -33,6 +34,7 @@ import org.w3c.dom.Node;
public abstract class SparkBaseTransformHandler implements
ISparkPipelineTransformHandler {
private SparkTransformMetricsAccumulator metricsAccumulator;
+ private SparkNativeMetricsListener metricsListener;
public void setMetricsAccumulator(SparkTransformMetricsAccumulator
metricsAccumulator) {
this.metricsAccumulator = metricsAccumulator;
@@ -42,6 +44,10 @@ public abstract class SparkBaseTransformHandler implements
ISparkPipelineTransfo
return metricsAccumulator;
}
+ public void setMetricsListener(SparkNativeMetricsListener metricsListener) {
+ this.metricsListener = metricsListener;
+ }
+
/**
* Attach native row tracking for this transform when a metrics accumulator
is registered. Returns
* {@code dataset} unchanged when metrics are disabled.
@@ -51,7 +57,7 @@ public abstract class SparkBaseTransformHandler implements
ISparkPipelineTransfo
if (dataset == null || transformMeta == null) {
return dataset;
}
- return SparkNativeMetrics.track(dataset, transformMeta.getName(),
metricsAccumulator, role);
+ return SparkNativeMetrics.track(dataset, transformMeta.getName(),
metricsListener, role);
}
@Override
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkFileInputHandler.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkFileInputHandler.java
index 851e57c27b..c21bec66ea 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkFileInputHandler.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/handler/SparkFileInputHandler.java
@@ -51,7 +51,11 @@ import org.apache.spark.sql.DataFrameReader;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
+import org.apache.spark.sql.types.DataType;
import org.apache.spark.sql.types.DataTypes;
+import org.apache.spark.sql.types.Metadata;
+import org.apache.spark.sql.types.StructField;
+import org.apache.spark.sql.types.StructType;
/**
* Native Spark file read for {@link SparkFileInputMeta}.
@@ -159,7 +163,19 @@ public class SparkFileInputHandler extends
SparkBaseTransformHandler {
boolean hasFields = meta.getFields() != null &&
!meta.getFields().isEmpty();
boolean nameBasedProjection = delimited && hasFields &&
!meta.isInferSchema();
- // For delimited + explicit fields: do NOT push StructType into the reader
(positional).
+ // For delimited + explicit fields: a reader StructType is bound by
POSITION, while the Hop
+ // field list is by name (subset, any order, case-insensitive). So the
schema we push is built
+ // from the file's own header: every file column in file order, typed
where a Hop numeric
+ // field matches, string otherwise. Same by-name projection/cast below
either way; numeric
+ // columns are then parsed once by the CSV converter instead of tokenized,
wrapped and cast,
+ // and no header-inference job is needed. Any problem reading the header →
schema-less load.
+ if (nameBasedProjection && meta.isHeader() && !meta.isMultiLine()) {
+ StructType typed =
+ typedSchemaFromHeader(log, spark, path, options, meta.getFields(),
transformMeta);
+ if (typed != null) {
+ reader = reader.schema(typed);
+ }
+ }
// Infer schema only when requested and no field list.
if (delimited && meta.isInferSchema() && !hasFields) {
reader = reader.option("inferSchema", "true");
@@ -330,18 +346,17 @@ public class SparkFileInputHandler extends
SparkBaseTransformHandler {
private static Column castColumn(Dataset<Row> source, String columnName,
SparkField field)
throws HopException {
Column c = col(columnName);
- // Always trim strings before numeric/date conversion
- Column trimmed = trim(c.cast(DataTypes.StringType));
+ int typeId = hopTypeId(field);
- String hopType = StringUtils.defaultIfBlank(field.getHopType(), "String");
- int typeId;
- try {
- typeId =
org.apache.hop.core.row.value.ValueMetaFactory.getIdForValueMeta(hopType);
- } catch (Exception e) {
- throw new HopException(
- "Unknown Hop type '" + hopType + "' for field '" + field.getName() +
"'", e);
+ // Already typed by the reader schema: no string round trip
+ DataType readerType = readerType(source, columnName);
+ if (readerType != null && readerType.equals(typedReaderType(typeId))) {
+ return c.alias(field.getName());
}
+ // Always trim strings before numeric/date conversion
+ Column trimmed = trim(c.cast(DataTypes.StringType));
+
return switch (typeId) {
case IValueMeta.TYPE_STRING, IValueMeta.TYPE_INET ->
trimmed.alias(field.getName());
case IValueMeta.TYPE_INTEGER -> trimmed.cast(DataTypes.LongType);
@@ -355,6 +370,174 @@ public class SparkFileInputHandler extends
SparkBaseTransformHandler {
};
}
+ private static int hopTypeId(SparkField field) throws HopException {
+ String hopType = StringUtils.defaultIfBlank(field.getHopType(), "String");
+ try {
+ return
org.apache.hop.core.row.value.ValueMetaFactory.getIdForValueMeta(hopType);
+ } catch (Exception e) {
+ throw new HopException(
+ "Unknown Hop type '" + hopType + "' for field '" + field.getName() +
"'", e);
+ }
+ }
+
+ private static DataType readerType(Dataset<Row> source, String columnName) {
+ for (StructField f : source.schema().fields()) {
+ if (f.name().equals(columnName)) {
+ return f.dataType();
+ }
+ }
+ return null;
+ }
+
+ /**
+ * Spark type the CSV converter may parse directly for a Hop type; null
keeps the column as string
+ * for the lenient cast path (dates need the format mask, booleans accept
Y/N, decimals need
+ * precision).
+ */
+ static DataType typedReaderType(int hopTypeId) {
+ return switch (hopTypeId) {
+ case IValueMeta.TYPE_INTEGER -> DataTypes.LongType;
+ case IValueMeta.TYPE_NUMBER -> DataTypes.DoubleType;
+ default -> null;
+ };
+ }
+
+ /**
+ * Build the reader schema from the first line of the (first) file: file
columns in file order,
+ * typed where a Hop numeric field matches the header name. Returns null
when the header cannot be
+ * read or is ambiguous, in which case the caller keeps the schema-less load.
+ */
+ static StructType typedSchemaFromHeader(
+ ILogChannel log,
+ SparkSession spark,
+ String path,
+ Map<String, String> options,
+ List<SparkField> fields,
+ TransformMeta transformMeta) {
+ try {
+ org.apache.hadoop.conf.Configuration conf =
spark.sessionState().newHadoopConf();
+ org.apache.hadoop.fs.Path p = new org.apache.hadoop.fs.Path(path);
+ org.apache.hadoop.fs.FileSystem fs = p.getFileSystem(conf);
+ org.apache.hadoop.fs.Path file = firstDataFile(fs, p);
+ if (file == null) {
+ return null;
+ }
+ String headerLine = readFirstLine(fs, conf, file,
options.getOrDefault("encoding", "UTF-8"));
+ if (StringUtils.isBlank(headerLine)) {
+ return null;
+ }
+ String[] names = parseHeader(headerLine, options);
+ if (names == null || names.length == 0) {
+ return null;
+ }
+ Map<String, Integer> typeByLower = new java.util.HashMap<>();
+ for (SparkField field : fields) {
+ if (StringUtils.isNotEmpty(field.getName())) {
+ typeByLower.put(field.getName().toLowerCase(Locale.ROOT),
hopTypeId(field));
+ }
+ }
+ Set<String> seen = new HashSet<>();
+ StructField[] structFields = new StructField[names.length];
+ for (int i = 0; i < names.length; i++) {
+ String name = names[i] == null ? "" : names[i];
+ if (name.isEmpty() || !seen.add(name)) {
+ // Spark would rename empty/duplicate headers; keep the proven path
instead
+ return null;
+ }
+ Integer typeId = typeByLower.get(name.toLowerCase(Locale.ROOT));
+ DataType type = typeId == null ? null : typedReaderType(typeId);
+ structFields[i] =
+ new StructField(
+ name, type == null ? DataTypes.StringType : type, true,
Metadata.empty());
+ }
+ return new StructType(structFields);
+ } catch (Exception e) {
+ if (log != null) {
+ log.logDetailed(
+ "Typed CSV schema not applied for '"
+ + transformMeta.getName()
+ + "' (falling back to schema-less load): "
+ + e.getMessage());
+ }
+ return null;
+ }
+ }
+
+ private static org.apache.hadoop.fs.Path firstDataFile(
+ org.apache.hadoop.fs.FileSystem fs, org.apache.hadoop.fs.Path p) throws
java.io.IOException {
+ org.apache.hadoop.fs.FileStatus[] matches = fs.globStatus(p);
+ if (matches == null || matches.length == 0) {
+ return null;
+ }
+ List<org.apache.hadoop.fs.FileStatus> candidates = new ArrayList<>();
+ for (org.apache.hadoop.fs.FileStatus status : matches) {
+ if (status.isDirectory()) {
+ for (org.apache.hadoop.fs.FileStatus child :
fs.listStatus(status.getPath())) {
+ if (child.isFile() && !isHidden(child.getPath())) {
+ candidates.add(child);
+ }
+ }
+ } else if (status.isFile() && !isHidden(status.getPath())) {
+ candidates.add(status);
+ }
+ }
+ if (candidates.isEmpty()) {
+ return null;
+ }
+ // Same choice Spark's header inference makes: the first file in path order
+ candidates.sort(java.util.Comparator.comparing(f ->
f.getPath().toString()));
+ return candidates.get(0).getPath();
+ }
+
+ private static boolean isHidden(org.apache.hadoop.fs.Path path) {
+ String name = path.getName();
+ return name.startsWith("_") || name.startsWith(".");
+ }
+
+ private static String readFirstLine(
+ org.apache.hadoop.fs.FileSystem fs,
+ org.apache.hadoop.conf.Configuration conf,
+ org.apache.hadoop.fs.Path file,
+ String encoding)
+ throws java.io.IOException {
+ org.apache.hadoop.io.compress.CompressionCodec codec =
+ new
org.apache.hadoop.io.compress.CompressionCodecFactory(conf).getCodec(file);
+ java.io.InputStream in = fs.open(file);
+ if (codec != null) {
+ in = codec.createInputStream(in);
+ }
+ try (java.io.BufferedReader reader =
+ new java.io.BufferedReader(new java.io.InputStreamReader(in,
encoding))) {
+ String line = reader.readLine();
+ if (line != null && !line.isEmpty() && line.charAt(0) == '\uFEFF') {
+ line = line.substring(1);
+ }
+ return line;
+ }
+ }
+
+ /** Split the header with the same delimiter/quote/escape the Spark CSV
reader will use. */
+ private static String[] parseHeader(String line, Map<String, String>
options) {
+ com.univocity.parsers.csv.CsvParserSettings settings =
+ new com.univocity.parsers.csv.CsvParserSettings();
+ com.univocity.parsers.csv.CsvFormat format = settings.getFormat();
+ format.setDelimiter(options.getOrDefault("sep",
options.getOrDefault("delimiter", ",")));
+ String quote = options.getOrDefault("quote", "\"");
+ if (!quote.isEmpty()) {
+ format.setQuote(quote.charAt(0));
+ }
+ String escape = options.getOrDefault("escape", "\\");
+ if (!escape.isEmpty()) {
+ format.setQuoteEscape(escape.charAt(0));
+ }
+ settings.setIgnoreLeadingWhitespaces(
+ Boolean.parseBoolean(options.getOrDefault("ignoreLeadingWhiteSpace",
"true")));
+ settings.setIgnoreTrailingWhitespaces(
+ Boolean.parseBoolean(options.getOrDefault("ignoreTrailingWhiteSpace",
"true")));
+ settings.setMaxColumns(20000);
+ return new com.univocity.parsers.csv.CsvParser(settings).parseLine(line);
+ }
+
/**
* Parse date/timestamp strings. Uses the field format mask when set;
otherwise tries common Hop
* masks (including {@code yyyy/MM/dd}).
diff --git
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkNativeMetricsTest.java
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkNativeMetricsTest.java
index 8d86a95c8b..b80c0e66e3 100644
---
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkNativeMetricsTest.java
+++
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkNativeMetricsTest.java
@@ -17,8 +17,11 @@
package org.apache.hop.spark.core;
+import static org.apache.spark.sql.functions.col;
+import static org.apache.spark.sql.functions.count;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.util.ArrayList;
@@ -51,6 +54,7 @@ class SparkNativeMetricsTest {
.config("spark.ui.showConsoleProgress", "false")
.config("spark.metrics.staticSources.enabled", "false")
.config("spark.driver.host", "localhost")
+ .config("spark.sql.shuffle.partitions", "3")
.getOrCreate();
}
@@ -61,84 +65,151 @@ class SparkNativeMetricsTest {
}
}
- @Test
- void trackReportsInputRoleAcrossPartitions() {
+ private static Dataset<Row> ids(int n, int partitions) {
StructType schema =
new StructType(
new StructField[] {
DataTypes.createStructField("id", DataTypes.LongType, false),
});
List<Row> rows = new ArrayList<>();
- for (long i = 0; i < 40; i++) {
+ for (long i = 0; i < n; i++) {
rows.add(RowFactory.create(i));
}
- Dataset<Row> input = spark.createDataFrame(rows, schema).repartition(4);
+ return spark.createDataFrame(rows, schema).repartition(partitions);
+ }
+ private static SparkTransformMetricsAccumulator accumulator(String name) {
SparkTransformMetricsAccumulator acc = new
SparkTransformMetricsAccumulator();
- spark.sparkContext().register(acc, "native-metrics-input");
+ spark.sparkContext().register(acc, name);
+ return acc;
+ }
+
+ private static SparkNativeMetricsListener
listen(SparkTransformMetricsAccumulator acc) {
+ SparkNativeMetricsListener listener = new SparkNativeMetricsListener(acc);
+ listener.addTo(spark.sparkContext());
+ return listener;
+ }
+
+ private static Map<String, SparkTransformMetricSlice> settle(
+ SparkNativeMetricsListener listener, SparkTransformMetricsAccumulator
acc) {
+ listener.flush(spark.sparkContext(), 10_000L);
+ listener.removeFrom(spark.sparkContext());
+ return acc.value();
+ }
+
+ private static long sum(
+ Map<String, SparkTransformMetricSlice> slices,
+ String transform,
+ java.util.function.ToLongFunction<SparkTransformMetricSlice> f) {
+ return slices.values().stream()
+ .filter(s -> transform.equals(s.getTransformName()))
+ .mapToLong(f)
+ .sum();
+ }
+
+ @Test
+ void trackReportsInputRoleAcrossPartitionsWithoutRddBarrier() {
+ Dataset<Row> input = ids(40, 4);
+ SparkTransformMetricsAccumulator acc = accumulator("native-metrics-input");
+ SparkNativeMetricsListener listener = listen(acc);
Dataset<Row> tracked =
- SparkNativeMetrics.track(input, "file-in", acc,
SparkNativeMetrics.Role.INPUT);
+ SparkNativeMetrics.track(input, "file-in", listener,
SparkNativeMetrics.Role.INPUT);
+ // Anchor only: the lineage stays a Dataset plan (no LogicalRDD from
createDataFrame(rdd))
+
assertTrue(tracked.queryExecution().optimizedPlan().toString().contains("CollectMetrics"));
+
assertFalse(tracked.queryExecution().optimizedPlan().toString().contains("LogicalRDD"));
assertEquals(40L, tracked.count());
- Map<String, SparkTransformMetricSlice> slices = acc.value();
+ Map<String, SparkTransformMetricSlice> slices = settle(listener, acc);
assertFalse(slices.isEmpty());
-
- long totalInput = 0;
- long totalWritten = 0;
Set<Integer> copies = new HashSet<>();
for (SparkTransformMetricSlice slice : slices.values()) {
assertEquals("file-in", slice.getTransformName());
assertTrue(slice.isFinished());
assertTrue(slice.getStartTimeMs() > 0, "partition should record start
time");
assertTrue(slice.getEndTimeMs() >= slice.getStartTimeMs(), "end should
be >= start");
- totalInput += slice.getLinesInput();
- totalWritten += slice.getLinesWritten();
copies.add(slice.getCopyNr());
}
- assertEquals(40L, totalInput);
- assertEquals(40L, totalWritten);
+ assertEquals(40L, sum(slices, "file-in",
SparkTransformMetricSlice::getLinesInput));
+ assertEquals(40L, sum(slices, "file-in",
SparkTransformMetricSlice::getLinesWritten));
assertTrue(copies.size() >= 2, "expected multi-partition copies, got " +
copies);
}
@Test
void trackReportsOutputAndTransformRoles() {
- StructType schema =
- new StructType(
- new StructField[] {
- DataTypes.createStructField("v", DataTypes.StringType, false),
- });
- List<Row> rows =
- List.of(RowFactory.create("a"), RowFactory.create("b"),
RowFactory.create("c"));
- Dataset<Row> input = spark.createDataFrame(rows, schema).repartition(2);
+ Dataset<Row> input = ids(3, 2);
+ SparkTransformMetricsAccumulator acc = accumulator("native-metrics-roles");
+ SparkNativeMetricsListener listener = listen(acc);
- SparkTransformMetricsAccumulator outAcc = new
SparkTransformMetricsAccumulator();
- spark.sparkContext().register(outAcc, "native-metrics-output");
Dataset<Row> outTracked =
- SparkNativeMetrics.track(input, "file-out", outAcc,
SparkNativeMetrics.Role.OUTPUT);
+ SparkNativeMetrics.track(input, "file-out", listener,
SparkNativeMetrics.Role.OUTPUT);
assertEquals(3L, outTracked.count());
- long outSum =
-
outAcc.value().values().stream().mapToLong(SparkTransformMetricSlice::getLinesOutput).sum();
- long readSum =
-
outAcc.value().values().stream().mapToLong(SparkTransformMetricSlice::getLinesRead).sum();
- assertEquals(3L, outSum);
- assertEquals(3L, readSum);
-
- SparkTransformMetricsAccumulator txAcc = new
SparkTransformMetricsAccumulator();
- spark.sparkContext().register(txAcc, "native-metrics-tx");
Dataset<Row> txTracked =
- SparkNativeMetrics.track(input, "sort", txAcc,
SparkNativeMetrics.Role.TRANSFORM);
+ SparkNativeMetrics.track(input, "sort", listener,
SparkNativeMetrics.Role.TRANSFORM);
assertEquals(3L, txTracked.count());
- long written =
-
txAcc.value().values().stream().mapToLong(SparkTransformMetricSlice::getLinesWritten).sum();
- long read =
-
txAcc.value().values().stream().mapToLong(SparkTransformMetricSlice::getLinesRead).sum();
- assertEquals(3L, written);
- assertEquals(3L, read);
+
+ Map<String, SparkTransformMetricSlice> slices = settle(listener, acc);
+ assertEquals(3L, sum(slices, "file-out",
SparkTransformMetricSlice::getLinesOutput));
+ assertEquals(3L, sum(slices, "file-out",
SparkTransformMetricSlice::getLinesRead));
+ assertEquals(0L, sum(slices, "file-out",
SparkTransformMetricSlice::getLinesInput));
+ assertEquals(3L, sum(slices, "sort",
SparkTransformMetricSlice::getLinesWritten));
+ assertEquals(3L, sum(slices, "sort",
SparkTransformMetricSlice::getLinesRead));
+ }
+
+ @Test
+ void countsAfterShuffleAreNotDoubledBySortSampling() {
+ // group by → sort: range partitioning samples the aggregate once before
the real pass
+ Dataset<Row> input = ids(100, 4).withColumn("k", col("id").mod(5));
+ SparkTransformMetricsAccumulator acc =
accumulator("native-metrics-shuffle");
+ SparkNativeMetricsListener listener = listen(acc);
+
+ Dataset<Row> in =
+ SparkNativeMetrics.track(input, "in", listener,
SparkNativeMetrics.Role.INPUT);
+ Dataset<Row> grouped =
in.groupBy(col("k")).agg(count(col("id")).alias("n"));
+ grouped =
+ SparkNativeMetrics.track(grouped, "group", listener,
SparkNativeMetrics.Role.TRANSFORM);
+ Dataset<Row> sorted = grouped.orderBy(col("k"));
+ sorted = SparkNativeMetrics.track(sorted, "sort", listener,
SparkNativeMetrics.Role.TRANSFORM);
+ assertEquals(5L, sorted.count());
+
+ Map<String, SparkTransformMetricSlice> slices = settle(listener, acc);
+ slices
+ .values()
+ .forEach(
+ sl ->
+ System.out.println(
+ "SLICE "
+ + sl.getTransformName()
+ + " copy="
+ + sl.getCopyNr()
+ + " in="
+ + sl.getLinesInput()
+ + " written="
+ + sl.getLinesWritten()
+ + " start="
+ + sl.getStartTimeMs()
+ + " end="
+ + sl.getEndTimeMs()));
+ System.out.println(sorted.queryExecution().executedPlan().toString());
+ assertEquals(100L, sum(slices, "in",
SparkTransformMetricSlice::getLinesInput));
+ assertEquals(5L, sum(slices, "group",
SparkTransformMetricSlice::getLinesWritten));
+ assertEquals(5L, sum(slices, "sort",
SparkTransformMetricSlice::getLinesWritten));
+ }
+
+ @Test
+ void tokensAreUniquePerListener() {
+ SparkTransformMetricsAccumulator acc = new
SparkTransformMetricsAccumulator();
+ SparkNativeMetricsListener a = new SparkNativeMetricsListener(acc);
+ SparkNativeMetricsListener b = new SparkNativeMetricsListener(acc);
+ String t1 = a.register("x", SparkNativeMetrics.Role.TRANSFORM);
+ String t2 = a.register("x", SparkNativeMetrics.Role.TRANSFORM);
+ assertNotEquals(t1, t2);
+ assertNotEquals(t1, b.register("x", SparkNativeMetrics.Role.TRANSFORM));
+ assertTrue(t1.startsWith(SparkNativeMetricsListener.TOKEN_PREFIX));
}
@Test
- void trackIsNoOpWithoutAccumulator() {
+ void trackIsNoOpWithoutListener() {
StructType schema =
new StructType(
new StructField[] {
diff --git
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/handler/SparkFileIoHandlersTest.java
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/handler/SparkFileIoHandlersTest.java
index f879c6a185..aaba2e368d 100644
---
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/handler/SparkFileIoHandlersTest.java
+++
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/handler/SparkFileIoHandlersTest.java
@@ -34,6 +34,7 @@ import org.apache.hop.core.variables.Variables;
import org.apache.hop.metadata.serializer.memory.MemoryMetadataProvider;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.spark.core.SparkNativeMetricsListener;
import org.apache.hop.spark.core.SparkTransformMetricSlice;
import org.apache.hop.spark.core.SparkTransformMetricsAccumulator;
import org.apache.hop.spark.engines.SparkPipelineRunConfiguration;
@@ -98,10 +99,14 @@ class SparkFileIoHandlersTest {
SparkTransformMetricsAccumulator metrics = new
SparkTransformMetricsAccumulator();
spark.sparkContext().register(metrics, "file-io-metrics");
+ // Native handlers anchor observe() nodes; the listener turns Spark's task
metrics into slices
+ SparkNativeMetricsListener metricsListener = new
SparkNativeMetricsListener(metrics);
+ metricsListener.addTo(spark.sparkContext());
Map<String, Dataset<Row>> map = new HashMap<>();
SparkFileInputHandler inputHandler = new SparkFileInputHandler();
inputHandler.setMetricsAccumulator(metrics);
+ inputHandler.setMetricsListener(metricsListener);
inputHandler.handleTransform(
LogChannel.GENERAL,
new Variables(),
@@ -121,6 +126,7 @@ class SparkFileIoHandlersTest {
assertEquals(3, read.count());
assertEquals(2, read.columns().length);
+ metricsListener.flush(spark.sparkContext(), 10_000L);
long inputRows =
metrics.value().values().stream()
.filter(s -> "read".equals(s.getTransformName()))
@@ -142,6 +148,7 @@ class SparkFileIoHandlersTest {
SparkFileOutputHandler outputHandler = new SparkFileOutputHandler();
outputHandler.setMetricsAccumulator(metrics);
+ outputHandler.setMetricsListener(metricsListener);
outputHandler.handleTransform(
LogChannel.GENERAL,
new Variables(),
@@ -161,6 +168,7 @@ class SparkFileIoHandlersTest {
assertEquals(0, map.get("write").count());
assertTrue(Files.exists(outDir));
+ metricsListener.flush(spark.sparkContext(), 10_000L);
long outputRows =
metrics.value().values().stream()
.filter(s -> "write".equals(s.getTransformName()))