Skip to content

feat(event-logs): coverage for partial, rolled, and retried logs - #22

Merged
jeffbrennan merged 1 commit into
mainfrom
feature/event-log-coverage
Sep 13, 2026
Merged

jeffbrennan merged 1 commit into
mainfrom
feature/event-log-coverage

Conversation

@jeffbrennan

Copy link
Copy Markdown
Owner

Implements brief 04, increments 1-3.

Increment 1 — plan fallback and partial parsing

parse_source() streams events and requires neither ApplicationStart nor an adaptive
execution update. Plans are ranked final adaptive > in-progress adaptive > the plan on
SQLExecutionStart, so a non-AQE workload no longer yields "No queries found" — the
nested_loop_join fixture goes from 1 query to 4.

  • 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; strict=True raises
    with URI and line number.
  • Missing data stays missing: no Task Metrics → metrics=None, no SQLExecutionEnd →
    null end timestamp and duration, no JobEnd → null job duration.
  • Unknown events are counted, not fatal. Empty logs produce typed empty frames.

Increment 2 — identity and attempts

  • Stage attempts are keyed by (stage_id, stage_attempt_id); task_status,
    task_failure_reason, partition_id, stage_status/num_tasks/failure_reason reach
    combined.
  • Job↔stage and query↔stage are association tables (dfs.job_stage, dfs.query_stage).
    Joining them into the task frame was duplicating rows: nested_final_plans reported 31
    rows for 19 tasks, counting 12 tasks' bytes twice.
  • Output accounting counts retained outputs — one attempt per (stage_id, partition),
    resolving losing speculative copies and partitions recomputed in a later stage attempt —
    while resource accounting counts every attempt. to_plan_summary()["total_basis"]
    records which basis each total used.
  • Tasks no SQL execution claims (schema inference, RDD work) keep a null query_id
    instead of being dropped.

Increment 3 — source discovery and streaming

  • New sparkparse/eventlog.py: discover_sources() collapses rolled eventlog_v2_*
    directories into one ordered source, flags .inprogress, and ignores markers and
    checksums. resolve_source() handles explicit selection (name, app id, rolling dir,
    path) with a documented newest-mtime default.
  • New sparkparse logs command and get --all-apps; capture can now select rolled and
    compressed logs.
  • Readable codecs: none, gz, zstd (new sparkparse[zstd] extra). Spark's lz4,
    lzf and snappy use Java-specific framing and fail with an error naming the codec.
  • Reads stream line by line everywhere, including the dashboard's log listing. A test
    asserts the parser never calls read()/readlines() on a log handle, and
    tests/benchmark_ingestion.py records elapsed time and process peak RSS per fixture
    size (measured: ~36 KB retained per task; peak RSS tracks retained models, not bytes
    read).

Notes for review

  • The golden fixtures are regenerated. Three things move: more queries parse, duplicate
    task rows are gone, and accumulator list ordering is now deterministically tie-broken on
    stage/task/accumulator id.
  • combined gained nine columns; anything reading it positionally needs updating.
  • Live validation of rolled/compressed logs on Databricks is still outstanding and is
    noted in the brief's status line.

Test plan

  • just ci — 421 passed, 2 skipped; ruff, format and pyrefly clean.
  • pytest tests/test_capture.py (local JVM) — 4 passed.
  • pytest tests/test.py (local JVM) — 10 passed, 1 skipped, 2 failed
    (test_python_udf, test_dpp_query); both reproduce identically on unmodified main,
    so they are environment issues, not regressions.
  • CLI smoke on data/logs/raw: logs, get, get --all-apps, analyze.

Implements brief 04, increments 1-3.

Plan fallback and partial parsing: parse_source() streams events and requires
neither ApplicationStart nor an adaptive update. Plans rank final adaptive >
in-progress adaptive > the plan on SQLExecutionStart, so a non-AQE workload no
longer yields "No queries found". A bad line is corrupt_line when more lines
follow and truncated_tail when it is last; strict mode raises with URI and line.
Missing task metrics stay None, unfinished queries and jobs keep null end
timestamps, unknown events are counted.

Identity and attempts: stage attempts are keyed by (stage_id,
stage_attempt_id); task status, failure reason and partition id reach the task
frame. Job/stage and query/stage relations move to association tables so they
cannot duplicate task rows -- nested_final_plans was reporting 31 rows for 19
tasks, double-counting 12. Output totals count retained outputs (one attempt per
partition, resolving speculative races and recomputed partitions); resource
totals count every attempt.

Source discovery and streaming: new sparkparse/eventlog.py discovers logical
sources, collapsing rolled eventlog_v2_* directories, flagging .inprogress and
ignoring markers. Explicit selection by name, app id, rolling dir or path, with
a documented newest-mtime default; new `sparkparse logs` and `get --all-apps`.
Readable codecs are none, gz and zstd; Spark's Java-framed lz4/lzf/snappy fail
by name. Reads never copy a log whole, asserted directly and benchmarked by
tests/benchmark_ingestion.py.

Golden fixtures are regenerated: more queries parse, duplicate task rows are
gone, and accumulator ordering is now deterministic.
@jeffbrennan
jeffbrennan merged commit 1092f37 into main Sep 13, 2026
3 checks passed
@jeffbrennan
jeffbrennan deleted the feature/event-log-coverage branch September 13, 2026 00:25
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant