sparkforensics-operator
An Airflow operator that runs sparkforensics analysis on a Spark job's event log right after it finishes, then acts on the result: persist the report, push it to XCom, notify (via any messaging/paging system you implement), and optionally fail the DAG if a budget was breached.
Airflow tells you a Spark task succeeded. It doesn't tell you that stage 14 spilled 40GB from a skewed join, or that runtime crept 20% past last week's run. This package closes that gap without a separate monitoring job: it runs as part of the DAG, right where the Spark task just ran.
Why this exists
- Two ways to wire it in: a standalone
SparkForensicsOperatorthat fails the DAG on a threshold breach, or aspark_forensics_callbackfactory to attach ason_success_callbackon the Spark task itself, log-only. - Where the log comes from and where the analysis runs are separate, composable choices. Fetch the log to the worker (Spark History Server, filesystem/mounted-HDFS path template, XCom, SFTP, or a History Server behind an SSH tunnel), or leave it where it is and point the analysis at it (a path on an SSH host, or a History Server application the CLI fetches itself).
- Run the analysis on the Airflow worker, or on an SSH host that already
sees the logs, such as the node running the Spark History Server. With
remote analysis the log never crosses the network, and the workers
(including managed ones like MWAA) don't need Node.js. With
deferrable=Truethe task also gives up its worker slot while the remote analysis runs, and a triggerer waits on it instead. - Budget thresholds mirror
sparkforensics-analyze's own CLI flags (max runtime, spill, skew, failed-task rate, min efficiency); only the ones you configure are enforced, and a breach raisesThresholdBreachedwithout consuming the task's retries, since re-running would just reach the same verdict. - Reports persist to a local path,
file://, ors3://, and a clickable "SparkForensics report" link shows up on the task in the Airflow UI. Setreport_url_template(an S3 console URL, say) to make it open in a browser. - A small summary (finding counts per impact band, breached thresholds,
the destination) goes to XCom under
sparkforensics_summary, so downstream tasks can branch without reading the report. - Hook arguments are Jinja templates (
path_template="/logs/{{ run_id }}",app_id="{{ ti.xcom_pull(...) }}"), rendered per task like any templated operator field. - Installs as an Airflow provider, listed by
airflow providers list. - An optional
Notifieryou implement (Slack, MS Teams, email, PagerDuty, ZenDuty, whatever you use) gets a best-effort pass/fail summary; a delivery failure never fails the task itself. - Tested against both Airflow 2.6+ and Airflow 3.0+ (CI matrix, Python 3.9-3.12). No live Spark, Airflow or Node is needed to run the tests: everything but the shell on the stand-in SSH host is mocked.
Quick start
pip install sparkforensics-operator
Standalone operator, chained after the Spark task, fails the DAG on breach:
from sparkforensics_operator import FilesystemLogSourceHook, SparkForensicsOperator, SubprocessAnalyzeHook
run_spark_job = SparkSubmitOperator(task_id="run_spark_job", ...)
check_spark_job = SparkForensicsOperator(
task_id="check_spark_job",
log_source=FilesystemLogSourceHook(path_template="/mnt/spark-logs/{{ run_id }}/eventlog"),
backend=SubprocessAnalyzeHook(),
report_dest="s3://reports/{{ run_id }}/report.json",
max_runtime_ms=3_600_000,
max_skew_ratio=3,
on_threshold_breach="fail",
)
run_spark_job >> check_spark_job
Or attach it as a callback instead. Airflow logs and swallows any exception
raised inside a callback, so a breach here never fails or retries the
upstream task, whatever on_threshold_breach is set to:
from sparkforensics_operator import SubprocessAnalyzeHook, XComLogSourceHook, spark_forensics_callback
run_spark_job = SparkSubmitOperator(
task_id="run_spark_job",
on_success_callback=spark_forensics_callback(
log_source=XComLogSourceHook(task_id="run_spark_job"),
backend=SubprocessAnalyzeHook(),
report_dest="s3://reports/run_spark_job/app.json",
max_runtime_ms=3_600_000,
),
...,
)
Or run the analysis on the Spark History Server's own node over SSH, so the
event log never reaches the worker. ssh_conn_id is the Airflow SSH
connection to that node, and base_url is the History Server as seen from
it:
from sparkforensics_operator import HistoryServerAppLogSourceHook, SparkForensicsOperator, SSHAnalyzeHook
check_spark_job = SparkForensicsOperator(
task_id="check_spark_job",
log_source=HistoryServerAppLogSourceHook(
base_url="http://localhost:18080",
app_id="{{ ti.xcom_pull(task_ids='run_spark_job', key='app_id') }}",
),
backend=SSHAnalyzeHook(ssh_conn_id="shs_node"),
report_dest="s3://reports/{{ run_id }}/report.json",
max_runtime_ms=3_600_000,
)
RemotePathLogSourceHook(ssh_conn_id="shs_node", path_template="/spark-logs/{{ run_id }}")
points SSHAnalyzeHook at an event log path on that host instead.
Add deferrable=True to that operator to release the worker slot while
the analysis runs: the CLI runs as a detached job on the SSH host and the
triggerer polls it, which needs a running triggerer with
sparkforensics-operator[ssh] installed and
apache-airflow-providers-ssh>=6.0.1, so Airflow 2.11 or newer. See
Deferrable execution.
Whichever host runs the analysis (the worker for SubprocessAnalyzeHook,
the SSH host for SSHAnalyzeHook) needs Node.js 18+ and the
sparkforensics-cli npm package (npm install -g sparkforensics-cli), so
sparkforensics-analyze is resolvable on PATH. Add the ssh extra for
the SSH hooks (it works from Airflow 2.6, like the package; only
deferrable=True needs Airflow 2.11+), and the s3 extra if
report_dest is s3://....
Full details in docs/runbook.md.
Development
pip install -e ".[test,s3,ssh]"
pytest
Test against a specific Airflow major version:
tox -e py311-airflow2
tox -e py311-airflow3
Learn more
- Architecture, components and design decisions
- API reference, the full operator/hook/threshold surface
- Runbook, worker and SSH host prerequisites, troubleshooting
- Glossary, domain terms
License
Release files for sparkforensics-operator 0.2.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| sparkforensics_operator-0.2.0.tar.gz | 109.5 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| sparkforensics_operator-0.2.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 164.4 kB
Release files / sparkforensics_operator-0.2.0.tar.gz
| Download URL | sparkforensics_operator-0.2.0.tar.gz |
|---|---|
| Size | 109.5 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
75d828e297ee13fbcd787bbd36aafd1b782149ace39966ee32e69e8809c99357
|
|
BLAKE2b-256 checksum How to use checksums |
be2997c028d3f5ad0e94983cbbdd1a5cddc55b393317a5b2a27686d88b72b98c
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 25, 2026.
Transparency logRelease files / sparkforensics_operator-0.2.0-py3-none-any.whl
| Download URL | sparkforensics_operator-0.2.0-py3-none-any.whl |
|---|---|
| Size | 54.9 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
e26683642ebaabff5f79d6a595ad89a68a0dbc03da4245ce5cd8e608400bb101
|
|
BLAKE2b-256 checksum How to use checksums |
886849ff1e374f5028a2c8709bcd0ef7d6002e1b16c5067d4efdaa4b37a9b51a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 25, 2026.
Transparency log