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()))

Reply via email to