mattcasters opened a new issue, #8648:
URL: https://github.com/apache/hop/issues/8648
### What would you like to happen?
### Background & Context
This proposal builds upon and synthesizes the ideas in:
- #6797: *A safe single threaded pipeline engine* (disk-spooling rowsets and
step resumption)
- #6537: *Dev Cache Mode / Step Checkpoints* (avoiding slow upstream
re-queries during pipeline development)
### The Real Problem: Operational Outages vs. Technical Retries
Distributed processing frameworks like Apache Spark or Apache Flink address
*technical infrastructure faults* (e.g., worker node failure, disk I/O errors,
rescheduling a failed task on a healthy executor). However, in real-world
production data engineering, the dominant causes of failure are *environmental
and operational*:
- An upstream source database undergoes maintenance or crashes for several
hours.
- A remote network connection drops.
- A cloud API rate-limits or revokes tokens.
In these operational failure scenarios, automated task retries fail
repeatedly. Worse, when the source database or API recovers hours later, its
state has mutated. If an extraction failed halfway through, an engine cannot
simply "continue where it left off" from the live database—doing so would
stitch together a temporally inconsistent "Frankenstein" dataset (e.g., 30%
from 02:00 AM and 70% from 05:00 AM).
Generic, automated resumption across arbitrary live data sources is
fundamentally flawed. Instead, we need to **manage the chaos** by giving data
engineers explicit tools to declare snapshot boundaries, contracts, and replay
policies at both the pipeline and workflow orchestration levels.
---
### Architectural Proposal
#### 1. Replay Gates (Sealing Gates) on Transforms and Actions
Instead of guessing where data boundaries lie, users explicitly mark
transforms or workflow actions as **Replay Gates**:
- **GUI Context Actions**: Modeled after `plugins/misc/debug`
(`TransformDebugGuiPlugin`), transforms and workflow actions provide a
`@GuiContextAction`: *"Mark as Replay Gate (Seal Output)"*.
- **Canvas Visual Decorator**: Modeled after
`DrawTransformDebugLevelBeeExtensionPoint`, using `PipelinePainterTransform`
and `WorkflowPainterAction` extension points to render an overlay badge (e.g. a
gate/seal icon) over marked components.
- **Attributes Map Persistence**: Settings are persisted cleanly in
`pipelineMeta.getAttributesMap().get("replay-gate")` and
`workflowMeta.getAttributesMap().get("replay-gate")`, keeping `.hpl` and `.hwf`
definitions clean and backward-compatible.
#### 2. Persistent Spooling RowSet & Resilient VFS Storage
- Implement `PersistentSpoolingRowSet` backed by Apache Commons VFS
(`HopVfs`).
- When a pipeline hits a Replay Gate, rows are serialized
(`rowMeta.writeData()`) with block compression (Snappy / ZSTD) to a designated
VFS path.
- **VFS Everywhere**: Spool paths can target local scratch (`file:///...`),
Amazon S3 (`s3://...`), Apache Hadoop HDFS (`hdfs://...`), or Azure Data Lake
(`azure://...`). If a Kubernetes pod or batch VM terminates mid-pipeline, the
spooled data and tamper-evident manifest survive in S3/HDFS.
- **Snapshot Manifest (`SnapshotManifest.json`)**: When the extraction
reaches completion without error, the gate writes a manifest containing the
`snapshotId`, row count, serialized `IRowMeta` schema, and SHA-256 checksum,
marking the snapshot status as `SEALED`. If interrupted, it is marked
`INTERRUPTED`.
#### 3. Pre-Flight Replay Evaluator (Safety Gate)
Before allowing an operator to replay a failed workflow, Hop runs a
deterministic pre-flight evaluation:
- **Upstream Gate Verification**: Verifies that all required Replay Gates
prior to the failure point exist, are marked `SEALED`, and pass SHA-256 / row
count validation.
- **Interruption Guard**: If a source extract was interrupted mid-stream,
downstream replay is **strictly blocked**, alerting the operator that the
source data was partially extracted and must be re-staged from scratch.
- **Divergence Detection**: Ensures pipeline schemas and workflow variables
have not incompatibly drifted.
#### 4. The Replaying Workflow Engine (`ReplayWorkflowEngine`)
A dedicated workflow engine plugin (`@WorkflowEnginePlugin(id = "Replay")`):
- Reads the previous execution record from the active
`IExecutionInfoLocation`.
- Bypasses verified, completed upstream actions while populating their
recorded `ActionResult` and workflow variables.
- Resumes execution directly at the designated action (or the first failed
action).
- For child pipelines, binds the remote VFS spooled snapshot so the pipeline
runs against the sealed data without touching the source database.
#### 5. Operations Experience in the Execution Perspective
In the Hop GUI **Execution Perspective** (`WorkflowExecutionViewer`):
- A new toolbar action: **"Replay Workflow..."**.
- Opens a Replay Dialog showing:
- A pre-flight checklist with status badges (green checkmarks for verified
sealed snapshots, red warnings for interrupted extracts).
- Selection of replay start point (e.g. *"Resume from first failed
action"*, *"Re-run entire workflow"*, or custom selection).
- Parameter / variable review and overrides.
- Initiates the replay run and links it to the original failed execution for
complete operational auditability.
### Issue Priority
Priority: 2
### Issue Component
Component: Workflows
--
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]