Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 50 additions & 3 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@ Dash dashboard for identifying performance bottlenecks.
Spark event log (JSONL)
│
▼
sparkparse/eventlog.py – discover logical log sources (rolled, compressed,
│ in-progress) and stream their events line by line
▼
sparkparse/parse.py – parse raw log into ParsedLog (Pydantic model)
│
▼
Expand All @@ -36,7 +39,9 @@ Spark event log (JSONL)
| File | Purpose |
| ----------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `sparkparse/models.py` | All Pydantic models. `NodeType` enum (30+ values), `ParsedLog`, `ParsedLogDataFrames`, `NODE_TYPE_DETAIL_MAP` (node type → detail model class) |
| `sparkparse/parse.py` | `get_parsed_metrics()` is the main entry point. `parse_spark_ui_tree()` converts indented ASCII plans to node graphs. `parse_log()` orchestrates everything. |
| `sparkparse/parse.py` | `get_parsed_metrics()` is the main entry point (`get_all_parsed_metrics()` for every application in a directory). `parse_source()` parses one logical log incrementally; `parse_spark_ui_tree()` converts indented ASCII plans to node graphs. |
| `sparkparse/eventlog.py`| Event-log source discovery and streaming. `discover_sources()` collapses rolled `eventlog_v2_*` directories into one source and skips markers/checksums; `resolve_source()` implements explicit selection and the newest-source default; `iter_lines()` streams segments. Readable codecs: none, `zstd`, `gz`. |
| `sparkparse/schemas.py` | Canonical typed empty frame schemas (`DAG_SCHEMA`, `COMBINED_SCHEMA`) shared by the event-log and Connect paths. |
| `sparkparse/clean.py` | `log_to_dag_df()` and `log_to_combined_df()` produce the two output DataFrames. `get_readable_size()` and `get_readable_timing()` are Polars expression helpers. |
| `sparkparse/app.py` | Typer CLI. `get` → parses and writes output files. `viz` → launches dashboard. |
| `sparkparse/connect.py` | Spark Connect adapter. Intercepts client action boundaries and `_build_metrics`, attributes metrics per execution/thread, and builds the dag frame. `probe_connect_support()` reports the client surface; `SparkConnectCapture.from_plan_metrics()` replays recorded executions offline. |
Expand All @@ -56,14 +61,53 @@ Spark event log (JSONL)

### `combined` DataFrame columns (per task)

- `log_name`, `parsed_log_name`, `query_id`, `query_function`
- `log_name`, `parsed_log_name`, `query_id`, `query_function`, `query_count`
- `query/job/stage/task` start/end timestamps and duration_seconds
- `task_id`, `executor_id`, `nodes` (list of physical plan nodes for this task)
- `stage_attempt_id`, `stage_num_tasks`, `stage_status`, `stage_failure_reason`
- `task_id`, `partition_id`, `attempt`, `task_status`, `task_succeeded`,
`task_failure_reason`, `executor_id`, `nodes` (plan nodes this task fed)
- Executor metrics: `executor_run_time_seconds`, `executor_cpu_time_seconds`, `jvm_gc_time_seconds`, `peak_execution_memory_bytes`
- Input/output: `bytes_read`, `records_read`, `bytes_written`, `records_written`
- Shuffle: `shuffle_remote_bytes_read`, `shuffle_local_bytes_read`, `shuffle_bytes_written`
- Spill: `memory_bytes_spilled`, `disk_bytes_spilled`

### Attempt and association contract

- One row per task **attempt**, keyed by `(stage_id, stage_attempt_id, task_id)`.
A retried stage keeps both attempts; keying on `stage_id` alone would fold a
retry's metrics into the original.
- Resource usage (`executor_run_time_seconds`, spill, GC) counts every attempt;
output accounting (`bytes_read/written`, `records_*`, shuffle bytes) counts
only *retained* outputs. Success alone is not enough: a losing speculative
copy and a partition recomputed in a later stage attempt both end
successfully, so `analyze.retained_outputs()` keeps one attempt per
`(stage_id, partition)` — latest stage attempt, earliest finish within it.
`to_plan_summary()["total_basis"]` records which basis each total used.
- A stage can belong to several jobs and serve several queries. Those relations
live in `ParsedLogDataFrames.job_stage` / `.query_stage`; joining them into
`combined` would duplicate task rows and double-count their metrics. The task
frame carries the earliest attributed query plus `query_count`.
- Tasks no SQL execution claims (schema inference, RDD work) stay in `combined`
with a null `query_id` rather than being dropped.

### Event-log source contract

- `parse_source()` requires neither `SparkListenerApplicationStart` nor an
adaptive execution update. The plan comes from the strongest event seen:
final adaptive plan > in-progress adaptive plan > the plan on
`SQLExecutionStart`. A non-AQE workload is normal, not an error.
- A line that fails to decode is `corrupt_line` when more lines follow and
`truncated_tail` when it is the last line of the last segment. Tolerant mode
records both in `ParsedLog.diagnostics`; `strict=True` raises with URI and
line number.
- Missing data stays missing: a task with no `Task Metrics` has `metrics=None`,
a query with no end event has a null end timestamp and duration.
- Reads are incremental: segments stream line by line and the log text is never
copied whole (`tests/test_eventlog.py` asserts no `read()`/`readlines()` on a
log handle). Peak RSS still grows with retained model state — roughly 36 KB
per task on the recorded fixtures. `python -m tests.benchmark_ingestion
<log_dir> [log_file]` reports elapsed time and peak RSS for a single log.

### Metric and findings contract

- Raw operator metrics are normalized through `sparkparse/metrics.py`. Use
Expand Down Expand Up @@ -204,6 +248,9 @@ uv run pyrefly check sparkparse/ tests/ # type check

## Known quirks

- Spark's `lz4`, `lzf` and `snappy` event-log codecs use Java-specific block
framing that no Python codec reads; those raise `UnsupportedCodecError` naming
the codec. `zstd` needs the `zstandard` package (`sparkparse[zstd]`).
- `capture.py` borrows supplied SparkSessions and requires event logging to be enabled
before capture. Use `cap.spark` inside the context; opt into an owned session explicitly.
- `test.py` and `test_capture.py` in `tests/` are integration tests that spin up a local
Expand Down
26 changes: 25 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,18 @@ pip install sparkparse
### CLI

```bash
# parse logs and write output files
# list the event-log sources found in a directory
sparkparse logs ./logs

# parse the newest log and write output files
sparkparse get --log-dir ./logs --out-format parquet

# parse one application explicitly (file name, rolling-log dir, or app id)
sparkparse get ./logs --log-file eventlog_v2_app-20260912-0001

# parse every application in the directory
sparkparse get ./logs --all-apps

# launch the dashboard
sparkparse viz --log-dir ./logs

Expand Down Expand Up @@ -57,6 +66,21 @@ sparkparse should own a local session, opt in explicitly with
with unavailable task/stage telemetry called out in `result.capabilities` rather than
represented as zeroes.

### event logs

Rolled logs (`spark.eventLog.rolling.enabled=true`) are read as one logical source
with ordered segments; `.inprogress` logs are read as far as they go and reported as
incomplete. Segments are streamed line by line, never copied into memory whole.
Zstd-compressed logs need `sparkparse[zstd]`; Spark's `lz4`, `lzf` and `snappy`
codecs use Java-specific framing that Python cannot decode, and fail with an error
naming the codec.

Neither `SparkListenerApplicationStart` nor adaptive execution is required. A
non-AQE query keeps the plan from its `SQLExecutionStart` event, a query with no
end event keeps a null duration, and a truncated final line is reported as
truncation rather than corruption. `strict=True` turns those diagnostics into
errors.

Use `backend="classic"` or `backend="connect"` to override detection. For post-run
ingestion without a Spark session, use `backend="event_log", log_file="/path/to/log"`.
Borrowed captures select the current application's log; use an explicit `log_file`
Expand Down
3 changes: 2 additions & 1 deletion plans/improvements-2026-09/04-event-logs.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
# 04 — Event-log coverage and scalable ingestion

Status: proposed. Priority: P1, with non-AQE support an early correctness fix.
Status: implemented offline (increments 1-3); live rolled/compressed Databricks
validation outstanding. Priority: P1, with non-AQE support an early correctness fix.
Depends on 01's identity/partial-result contract. Files: `parse.py`, `clean.py`,
`models.py`, `storage.py` and event-log fixtures.

Expand Down
2 changes: 1 addition & 1 deletion plans/improvements-2026-09/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ access mode. Spark Connect is a transport, not proof that compute is serverless.
| P0 | [01 — Capture and capabilities](01-capture-and-capabilities.md) | Safe session lifecycle, common finalization, explicit missing data | Implemented + smoke-tested |
| P0 | [02 — Connect correctness](02-connect-correctness.md) | Reliable query attribution and tolerant operator handling | Implemented offline; live validation outstanding |
| P0 | [03 — Analysis correctness and depth](03-analysis.md) | Accurate metrics and evidence-based findings | Implemented offline; increments 1–3 |
| P1 | [04 — Event-log robustness](04-event-logs.md) | Non-AQE, partial, rolled, retried, and larger workloads | Large; 01 identities |
| P1 | [04 — Event-log robustness](04-event-logs.md) | Non-AQE, partial, rolled, retried, and larger workloads | Implemented offline; increments 1–3 |
| P1 | [05 — Developer experience and validation](05-developer-experience.md) | Installable CLI, reproducible checks, serverless-friendly reports | Medium; packaging can start immediately |
| P1 | [06 — History and comparisons](06-history.md) | Comparable runs and meaningful regression alerts | Medium; 01 and 03 |

Expand Down
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@ azure = ["adlfs>=2024.1.0"]
gcs = ["gcsfs>=2024.1.0"]
cloud = ["s3fs>=2024.1.0", "adlfs>=2024.1.0", "gcsfs>=2024.1.0"]
delta = ["deltalake>=0.25"]
# Zstd-compressed event logs. Spark's lz4/lzf/snappy codecs use Java-specific
# framing and cannot be read from Python at all.
zstd = ["zstandard>=0.22"]


[tool.hatch.build.targets.wheel]
Expand All @@ -50,6 +53,7 @@ dev-dependencies = [
"pyspark>=3.5.5",
"pytest>=8.3.5",
"ruff>=0.9.9",
"zstandard>=0.25.0",
]

[tool.ruff]
Expand Down
83 changes: 80 additions & 3 deletions sparkparse/analyze.py
Original file line number Diff line number Diff line change
Expand Up @@ -446,6 +446,78 @@ def expr(self, value: Any) -> Any:

_COMPACT_DEFAULT_NODES = 25

# Output that a failed or killed attempt produced is discarded and recomputed,
# so counting it would inflate the ledger. The resources that attempt burned
# were still spent, so those are counted for every attempt.
OUTPUT_ACCOUNTING_COLUMNS: frozenset[str] = frozenset(
{
"bytes_read",
"records_read",
"bytes_written",
"records_written",
"shuffle_bytes_read",
"shuffle_bytes_written",
}
)


def _partition_key(columns: set[str]) -> pl.Expr | None:
"""Expression identifying the partition a task attempt computed."""
available = [name for name in ("partition_id", "index") if name in columns]
if not available:
return None
return pl.coalesce([pl.col(name) for name in available]).alias("_partition")


def retained_outputs(combined: pl.DataFrame) -> pl.DataFrame:
"""Rows for the task attempts whose output was actually kept.

Success is necessary but not sufficient. Two successful attempts can exist
for the same partition — a speculative copy that finished after the commit
was already awarded, or a partition recomputed in a later stage attempt
after its output was lost — and only one of them contributes bytes and
rows. Counting both inflates every output total.

One attempt survives per ``(stage_id, partition)``: the latest stage
attempt (its output supersedes the lost one), and within it the attempt
that finished first, which is the one Spark's commit coordinator would have
authorized. Every attempt stays in ``combined`` for resource accounting.

Sources that do not report task status (Spark Connect plan metrics) have no
``task_succeeded`` column; every row they do report is treated as kept.
"""
columns = set(combined.columns)
if "task_succeeded" not in columns:
return combined

kept = combined.filter(pl.col("task_succeeded").fill_null(True))
partition = _partition_key(columns)
if partition is None or "stage_id" not in columns or kept.is_empty():
return kept

order = [("stage_id", False), ("_partition", False)]
if "stage_attempt_id" in columns:
order.append(("stage_attempt_id", True))
if "task_end_timestamp" in columns:
order.append(("task_end_timestamp", False))
if "task_id" in columns:
order.append(("task_id", False))

return (
kept.with_columns(partition)
.sort(
[name for name, _ in order],
descending=[descending for _, descending in order],
nulls_last=True,
)
.unique(subset=["stage_id", "_partition"], keep="first", maintain_order=True)
.drop("_partition")
)


def total_basis(column: str) -> str:
return "retained_outputs" if column in OUTPUT_ACCOUNTING_COLUMNS else "all_attempts"


def _safe_node_name(row: dict[str, Any], redactor: Redactor) -> str | None:
"""Return the node's display name, redacted when it carries workload text.
Expand Down Expand Up @@ -606,9 +678,11 @@ def to_plan_summary(
# An empty task table is a valid Connect result, not evidence of zero work.
agg: dict[str, Any] = dict.fromkeys(_TOTAL_COLUMNS)
else:
agg = combined.select(
*(pl.sum(column).alias(column) for column in _TOTAL_COLUMNS)
).row(0, named=True)
kept = retained_outputs(combined)
agg = {}
for column in _TOTAL_COLUMNS:
frame = kept if column in OUTPUT_ACCOUNTING_COLUMNS else combined
agg[column] = frame.select(pl.sum(column)).item() if frame.height else None

summary: dict[str, Any] = {
"schema_version": SUMMARY_SCHEMA_VERSION,
Expand All @@ -618,6 +692,9 @@ def to_plan_summary(
"total_units": {
column: _TOTAL_UNITS[column].value for column in _TOTAL_COLUMNS
},
# Which attempts each total covers: resource usage counts every
# attempt, output accounting counts only the attempts that were kept.
"total_basis": {column: total_basis(column) for column in _TOTAL_COLUMNS},
"coverage": capabilities.model_dump(mode="json"),
}

Expand Down
69 changes: 67 additions & 2 deletions sparkparse/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,9 @@
from sparkparse import alerts, history
from sparkparse.analyze import to_analysis_export, to_plan_summary
from sparkparse.dashboard import init_dashboard, run_app
from sparkparse.eventlog import discover_sources
from sparkparse.models import OutputFormat, ParsedLogDataFrames, RunRecord
from sparkparse.parse import get_parsed_metrics
from sparkparse.parse import get_all_parsed_metrics, get_parsed_metrics
from sparkparse.storage import (
get_path_name,
get_path_stem,
Expand Down Expand Up @@ -75,8 +76,19 @@ def get(
] = "data/logs/raw",
log_file: Annotated[
str | None,
typer.Option(help="Parse a single log file instead of the whole directory."),
typer.Option(
help="Event log to parse: file name, rolling-log directory, "
"application id, or full path. Default: the newest source in "
"log_dir (most recent mtime, name order when there is none)."
),
] = None,
all_apps: Annotated[
bool,
typer.Option(
"--all-apps",
help="Parse every application in log_dir instead of just the newest.",
),
] = False,
out_dir: Annotated[
str | None, typer.Option(help="Directory to write parsed output files.")
] = "data/logs/parsed",
Expand All @@ -94,6 +106,22 @@ def get(
] = False,
) -> ParsedLogDataFrames:
"""Parse Spark event logs and write structured DataFrames to disk."""
if all_apps:
if log_file is not None:
raise typer.BadParameter("--all-apps cannot be combined with --log-file")
results = get_all_parsed_metrics(
log_dir=log_dir,
out_dir=out_dir,
out_format=out_format,
verbose=verbose,
strict=strict,
)
# Query ids restart per application, so the results stay separate. The
# command returns the newest one for interactive use; every one of them
# is written to out_dir.
typer.echo(f"Parsed {len(results)} application(s): {', '.join(results)}")
return results[sorted(results)[-1]]

return get_parsed_metrics(
log_dir=log_dir,
log_file=log_file,
Expand All @@ -105,6 +133,43 @@ def get(
)


@app.command("logs")
def logs(
log_dir: Annotated[
str, typer.Argument(help="Directory containing raw Spark event logs.")
] = "data/logs/raw",
format: Annotated[
AnalysisFormat,
typer.Option(help="Output format: 'text' (default) or 'json'."),
] = AnalysisFormat.text,
) -> None:
"""List the event-log sources discovered in a directory.

Rolling logs collapse into one source with ordered segments; marker files
and checksums are ignored. The last row is the one a bare parse selects.
"""
sources = discover_sources(log_dir)
if not sources:
typer.echo(f"No event log sources found in {log_dir}")
raise typer.Exit(1)

ordered = sorted(sources, key=lambda item: (item.modified or 0, item.name))
if format == AnalysisFormat.json:
typer.echo(json.dumps([s.model_dump(mode="json") for s in ordered], indent=2))
return

for source in ordered:
codecs = ",".join(codec.value for codec in source.codecs)
typer.echo(
f"{source.name}\t"
f"segments={len(source.segments)}\t"
f"rolling={source.rolling}\t"
f"complete={source.complete}\t"
f"codec={codecs}"
)
typer.echo(f"\nNewest (default selection): {ordered[-1].name}", err=True)


@app.command("analyze")
def analyze(
log_dir: Annotated[
Expand Down
Loading
Loading