Repository navigation
feat(event-logs): coverage for partial, rolled, and retried logs - #22
Merged
Merged
Conversation
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Implements brief 04, increments 1-3.
Increment 1 — plan fallback and partial parsing
parse_source()streams events and requires neitherApplicationStartnor an adaptiveexecution update. Plans are ranked final adaptive > in-progress adaptive > the plan on
SQLExecutionStart, so a non-AQE workload no longer yields "No queries found" — thenested_loop_joinfixture goes from 1 query to 4.corrupt_linewhen more lines follow andtruncated_tailwhen it is the last line of the last segment;strict=Trueraiseswith URI and line number.
Task Metrics→metrics=None, noSQLExecutionEnd→null end timestamp and duration, no
JobEnd→ null job duration.Increment 2 — identity and attempts
(stage_id, stage_attempt_id);task_status,task_failure_reason,partition_id,stage_status/num_tasks/failure_reasonreachcombined.dfs.job_stage,dfs.query_stage).Joining them into the task frame was duplicating rows:
nested_final_plansreported 31rows for 19 tasks, counting 12 tasks' bytes twice.
(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.
query_idinstead of being dropped.
Increment 3 — source discovery and streaming
sparkparse/eventlog.py:discover_sources()collapses rolledeventlog_v2_*directories into one ordered source, flags
.inprogress, and ignores markers andchecksums.
resolve_source()handles explicit selection (name, app id, rolling dir,path) with a documented newest-mtime default.
sparkparse logscommand andget --all-apps; capture can now select rolled andcompressed logs.
gz,zstd(newsparkparse[zstd]extra). Spark'slz4,lzfandsnappyuse Java-specific framing and fail with an error naming the codec.asserts the parser never calls
read()/readlines()on a log handle, andtests/benchmark_ingestion.pyrecords elapsed time and process peak RSS per fixturesize (measured: ~36 KB retained per task; peak RSS tracks retained models, not bytes
read).
Notes for review
task rows are gone, and accumulator list ordering is now deterministically tie-broken on
stage/task/accumulator id.
combinedgained nine columns; anything reading it positionally needs updating.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 unmodifiedmain,so they are environment issues, not regressions.
data/logs/raw:logs,get,get --all-apps,analyze.