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]