mattcasters commented on issue #6797:
URL: https://github.com/apache/hop/issues/6797#issuecomment-5854900734

   ### Architectural Examination: Safe Single-Threaded Engine (#6797), Dev 
Checkpoints (#6537), Spark Comparisons, and NiFi-Style Replay
   
   Here is a structured analysis exploring how a safe single-threaded engine 
could work, how it intersects with developer checkpoints, why Apache Spark does 
not address these use cases, and the inherent complexities of introducing 
NiFi-style "replay" into Hop.
   
   ---
   
   ### 1. How the Idea Behind Issue #6797 Could Work
   *(A safe, single-threaded pipeline engine with disk persistence and 
resumption)*
   
   The core premise of #6797 is to execute transforms **sequentially 
one-by-one** (never in parallel threads or multiple copies), spool intermediate 
rowsets to disk in a configurable location, and automatically resume execution 
from the failed transform upon restart.
   
   #### Engine Baseline
   Hop already has `SingleThreadedPipelineEngine` and 
`SingleThreadedPipelineExecutor`, but they historically operated in 
micro-batches or assumed row buffers could reside in memory (which deadlocks or 
runs out of memory if an upstream transform produces more rows than fit in an 
in-memory `IRowSet`). Hop also recently introduced `SpillingRowSet`, which 
spills excess rows to temporary files using `HopVfs` and `rowMeta.writeData()`, 
but it deletes segments as soon as they are consumed.
   
   #### Architecture for #6797
   To realize #6797:
   1. **Topological Staged Execution**:
      - The engine computes the topological order of the pipeline DAG.
      - Instead of scheduling all transforms concurrently on separate threads, 
**Transform 1** executes exclusively until it exhausts its source and calls 
`setDone()`.
      - All output rows are spooled to a persistent disk buffer.
      - Once Transform 1 finishes and its output files are flushed, **Transform 
2** initializes, reads from the spooled file, processes all rows, and spools to 
its own output file.
      - Multi-input transforms (e.g., `Stream Lookup`, `Merge Join`) wait until 
all upstream dependency transforms reach completion before beginning execution.
   
   2. **Persistent Spooling RowSet (`PersistentSpoolingRowSet`)**:
      - Unlike `SpillingRowSet`, this implementation retains spooled files 
after consumption until the entire pipeline run is complete or explicitly 
cleared.
      - Utilizes `HopVfs` to support configurable spool paths (local NVMe, 
scratch mounts, or cloud volumes).
      - Format:
        - Header: Serialized `IRowMeta` (schema, data types, field names, 
origins).
        - Data: Binary row serialization (`rowMeta.writeData(dos, row)`) with 
block compression (e.g., Snappy or ZSTD) to mitigate disk I/O bottlenecks.
        - Multiple downstream consumers (fan-out / copy hops) attach 
independent file read pointers to the same spooled segment.
   
   3. **Execution Manifest & Resumption Protocol**:
      - The engine writes an `execution-journal.json` in the run directory 
capturing:
        - Pipeline metadata fingerprint (to verify the pipeline structure has 
not changed).
        - State of each transform (`PENDING`, `RUNNING`, `COMPLETED`, `FAILED`).
        - URIs of the spooled output files.
        - Row metrics (lines read, written, input, output, errors).
      - **Resumption**:
        - If Transform 4 fails halfway through:
          - Transforms 1, 2, and 3 are recorded as `COMPLETED`.
          - When the pipeline is re-run with the same run ID/directory, the 
engine validates the pipeline fingerprint and **skips Transforms 1, 2, and 3**.
          - Transform 4's input `IRowSet` connects directly to the completed 
spool file of Transform 3.
          - Transform 4's output spool is reset, and execution resumes from 
Transform 4.
   
   4. **Trade-offs & Considerations**:
      - **I/O Overhead**: Spooling every intermediate hop on large datasets can 
introduce significant disk I/O amplification. Compression and fast 
serialization are critical.
      - **Idempotency & External Side Effects**: If a transform that writes to 
a database or calls external APIs fails midway, re-running that transform 
requires either transactional rollback or idempotent target operations (e.g., 
upserts) to prevent duplicate writes.
   
   ---
   
   ### 2. Requirements Behind Issue #6537
   *(Dev Cache Mode / Step Checkpoints)*
   
   Where #6797 addresses batch safety and crash recovery in production/staging 
runs, **#6537** focuses on **developer velocity during design and debugging**. 
A slow upstream extract (e.g., a 5-minute database query or remote file fetch) 
currently forces the developer to re-run from the beginning whenever tweaking 
downstream formulas or regexes.
   
   #### Key Requirements:
   1. **Interactive Checkpoint Gestures**:
      - Right-click any transform on the canvas -> *"Set Dev Checkpoint"* (or 
*"Cache Step Output"*).
      - Visual indicator on the canvas (e.g., a green cache badge on the 
transform icon).
      - Context actions: *"Clear Checkpoint Cache"*, *"Refresh Cache on Next 
Run"*, and *"Run From Checkpoint"*.
   2. **Row Limits & Ephemeral Storage**:
      - Configurable row cap (e.g., 1,000–10,000 rows) so dev caches remain 
lightweight and fast.
      - Storage in local project scratch space (e.g., 
`${PROJECT_HOME}/.hop/cache/`), completely isolated from the `.hpl` file so Git 
repositories remain clean.
   3. **DAG Virtualization / Upstream Pruning at Runtime**:
      - When running from a checkpoint:
        - Upstream transforms prior to the checkpoint are not initialized (no 
database connections opened, no remote endpoints contacted).
        - Downstream transforms read from an internal cache reader feeding rows 
directly from the checkpoint store.
   4. **Cache Invalidation & Dirty Tracking**:
      - If upstream transform metadata or pipeline parameters change, the GUI 
warns the user that upstream logic has diverged and offers to refresh the cache.
   5. **Synergy with #6797**:
      - #6537 is effectively an interactive GUI frontend over the persistent 
spooling mechanism required for #6797. Both rely on a disk-backed, schema-aware 
row serializer.
   
   ---
   
   ### 3. Why Spark Isn't Enough for Developers & Ops Use Cases
   
   While Apache Spark provides DAG scheduling, RDD lineage, and 
`df.checkpoint()`, it is fundamentally geared toward large-scale distributed 
analytics rather than the developer and operational workflows where Hop excels:
   
   | Dimension | Apache Hop | Apache Spark |
   | :--- | :--- | :--- |
   | **Inner Loop Latency** | **Sub-second**. Instant data preview and 
interactive runs on sample rows in 200–500 ms directly in the GUI. | **10–30+ 
seconds**. Cold JVM startup, `SparkSession` initialization, Catalyst 
optimization, and DAG stage compilation break interactive developer flow. |
   | **Execution Transparency** | **Visual, step-by-step observability**. 
Real-time row throughput (`r/w/i/o/e`) and buffer levels visible on every hop. 
| **Black-box execution**. Catalyst collapses operations into whole-stage 
codegen. Stack traces come from generated bytecode 
(`GeneratedClass$SpecificUnsafeProjection`) detached from visual steps. |
   | **Operational ETL vs. Analytical Big Data** | Excels at complex 
operational integration: legacy formats, paginated REST APIs, stored 
procedures, SFTP polling, and transactional database lookups. | Built for bulk 
relational operations over columnar stores (Parquet/Iceberg). Very inefficient 
at high-concurrency row-by-row API calls or transactional OLTP updates. |
   | **Dead-Letter Error Routing** | **Per-row error handling**. Failed rows 
are diverted to dedicated error streams with error codes and descriptions while 
the pipeline continues processing. | **All-or-nothing task failures**. A single 
unhandled exception or data mismatch can fail an entire task/stage, requiring 
complex defensive coding. |
   | **Granular Resumption** | Can resume at a specific failed transform 
without re-querying sensitive production OLTP source systems. | Partition 
retries handle transient faults, but persistent pipeline failures typically 
require re-running entire jobs or relying on external snapshot/table 
orchestrations. |
   | **Infrastructure Footprint** | Lightweight single JVM running in a 512 MB 
– 2 GB container at the edge, in VMs, or on developer laptops. | Heavy cluster 
footprint requiring cluster managers (YARN, K8s), distributed file systems, and 
dedicated driver/executor allocations. |
   
   ---
   
   ### 4. Complexities of Introducing Apache NiFi "Replay" in a Hop Setting
   
   In Apache NiFi, the **Data Provenance Replay** feature allows an operator to 
inspect any past event (e.g., dropped FlowFiles or processor failures), click 
**"Replay"**, and NiFi re-injects that exact FlowFile into the processor’s 
incoming queue.
   
   Attempting to introduce this into Hop reveals fundamental architectural 
divergences:
   
   1. **Discrete FlowFiles vs. Tabular Row Streaming**:
      - **NiFi**: Operates on coarse-grained **FlowFiles** (e.g., a file, 
document, or discrete message) containing string attributes and a pointer to an 
immutable byte stream in the Content Repository. Replaying 1 or 100 FlowFiles 
is low-overhead.
      - **Hop**: Processes high-throughput tabular row tuples (`Object[]` 
arrays) streaming continuously. Generating provenance events and retaining data 
for millions of individual rows at every transform would cause massive storage 
and I/O amplification.
   2. **Immutable Content Claims vs. Ephemeral In-Memory Rows**:
      - **NiFi**: Content in the Content Repository is append-only and 
immutable. Replay simply creates a new FlowFile pointing to existing byte 
offsets.
      - **Hop**: Rows are created, mutated, and garbage-collected across 
in-memory 
[`IRowSet`](file:///home/matt/git/mattcasters/hop/core/src/main/java/org/apache/hop/core/IRowSet.java)
 queues. There is no central immutable content repository unless the pipeline 
runs on a dedicated spooling engine.
   3. **`IRowMeta` Schema Coupling**:
      - In NiFi, payload bytes are self-contained. In Hop, an `Object[]` row is 
untyped and positional; its meaning depends strictly on the `IRowMeta` at that 
exact transform hop. If a pipeline is modified (columns added, removed, or 
reordered), replaying previously captured rows will cause array index or type 
mismatches.
   4. **Multi-Stream Synchronization**:
      - Transforms like `Merge Join`, `Stream Lookup`, and `Sorted Merge` 
depend on precise coordination between multiple streams. If rows are replayed 
on one incoming stream, the other stream has already advanced or completed, 
completely breaking join alignment.
   5. **Opaque Internal Transform State (`ITransformData`)**:
      - Hop transforms maintain transient state (running totals, hash lookups, 
sort buffers, open database connections) in their internal `data` object. NiFi 
processors operate on discrete transactional `ProcessSession` boundaries. In 
Hop, there is no lifecycle hook to "rewind" or re-initialize transform state 
for partial replays.
   6. **Side Effects & Non-Idempotent Sinks**:
      - Hop transforms regularly commit database writes in batches (e.g., every 
1,000 or 10,000 rows). If a transform fails halfway through, earlier batches 
are already committed. Replaying upstream rows without idempotency or rollback 
causes primary key collisions or duplicate records.
   
   ---
   
   ### Recommended Path Forward for Hop
   
   Rather than attempting to mirror NiFi's per-event provenance replay (which 
conflicts with Hop’s streaming tabular architecture):
   
   1. **Unify #6797 and #6537 around a shared persistent spooling layer**:
      - Build a `PersistentSpoolingRowSet` backed by `HopVfs` and compressed 
binary row serialization.
      - **For #6537 (Dev Checkpoints)**: Use this spooler with a row cap (e.g., 
1,000–10,000 rows) and canvas actions to cache slow upstream outputs during 
pipeline development.
      - **For #6797 (Safe Single-Threaded Engine)**: Use this spooler with 
topological staged execution and an execution manifest to provide step-level 
failure resumption in production/batch environments.
   2. **Leverage Native Error Handling for Row-Level Replays**:
      - Hop already allows routing rejected rows to an error handling stream. 
Persisting these rejected rows into a standardized dead-letter store provides a 
clean, safe, and idempotent mechanism to inspect, fix, and re-inject failed 
records without the overhead of a full provenance replay engine.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to