Skip to content

feat(connect): reliable Spark Connect capture - #20

Merged
jeffbrennan merged 9 commits into
mainfrom
claude/brief-02-changes-4rzz77
Sep 12, 2026
Merged

jeffbrennan merged 9 commits into
mainfrom
claude/brief-02-changes-4rzz77

Conversation

@jeffbrennan

Copy link
Copy Markdown
Owner

Rework the Spark Connect adapter so query attribution, operator mapping, and
timing semantics are established rather than assumed (brief 02).

  • Intercept action boundaries (to_table, to_pandas, to_table_as_iterator,
    execute_command, execute_command_as_iterator) so each action is one execution
    with its own logical plan, metrics, and server operation id; metrics are
    attributed by calling thread, never by list position.
  • Probe the client surface before the workload runs: required methods, the
    _build_metrics signature, and DataFrame.executionInfo availability. Install
    hooks atomically and restore the exact originals (no leftover instance
    attributes) on setup failure, workload failure, and normal exit. Reject a
    second capture on one client.
  • Treat metric batches as snapshots (last value per plan id wins) instead of
    draining the _build_metrics generator, which previously emptied the caller's
    metrics as well.
  • Attach join keys only on a logical plan-id match, a single-join-per-query
    mapping, or an operator name that carries them; otherwise leave them
    unresolved with a diagnostic. Joins record input_roles: unordered.
  • Degrade unrecognized operators to NodeType.Unknown with raw names and graph
    edges preserved, unless strict=True is forwarded from the capture API.
  • Use client-observed elapsed time for query duration (capability capped at
    partial, reason states the transfer inclusion) and keep cumulative operator
    time per node with its source unit.
  • Surface Connect diagnostics in CaptureResult and carry the operation id in
    source_execution_id.
  • Add tests/test_connect.py (41 offline cases, no Spark/gRPC) and a synthetic
    sanitized replay fixture.

Co-Authored-By: Claude Opus 5 noreply@anthropic.com
Claude-Session: https://claude.ai/code/session_01AcyKZJKNnbkRiaa44qysbU

jeffbrennan and others added 9 commits September 12, 2026 21:07
Rework the Spark Connect adapter so query attribution, operator mapping, and
timing semantics are established rather than assumed (brief 02).

- Intercept action boundaries (to_table, to_pandas, to_table_as_iterator,
  execute_command, execute_command_as_iterator) so each action is one execution
  with its own logical plan, metrics, and server operation id; metrics are
  attributed by calling thread, never by list position.
- Probe the client surface before the workload runs: required methods, the
  _build_metrics signature, and DataFrame.executionInfo availability. Install
  hooks atomically and restore the exact originals (no leftover instance
  attributes) on setup failure, workload failure, and normal exit. Reject a
  second capture on one client.
- Treat metric batches as snapshots (last value per plan id wins) instead of
  draining the _build_metrics generator, which previously emptied the caller's
  metrics as well.
- Attach join keys only on a logical plan-id match, a single-join-per-query
  mapping, or an operator name that carries them; otherwise leave them
  unresolved with a diagnostic. Joins record input_roles: unordered.
- Degrade unrecognized operators to NodeType.Unknown with raw names and graph
  edges preserved, unless strict=True is forwarded from the capture API.
- Use client-observed elapsed time for query duration (capability capped at
  partial, reason states the transfer inclusion) and keep cumulative operator
  time per node with its source unit.
- Surface Connect diagnostics in CaptureResult and carry the operation id in
  source_execution_id.
- Add tests/test_connect.py (41 offline cases, no Spark/gRPC) and a synthetic
  sanitized replay fixture.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AcyKZJKNnbkRiaa44qysbU
- notebooks/validate_connect_capture.py: bounded serverless checks for brief 02
  (action coverage per intercepted client method, timing semantics, unknown
  operator degradation, two-join key isolation, hook lifecycle) that also prints
  a sanitized PlanMetrics fixture for offline replay. Uses spark.range frames and
  a noop sink only; no table is created, read, or modified.
- databricks.yml: sparkparse_validate_connect job for that notebook.
- SparkConnectCapture.to_plan_metrics(): export recorded executions in the shape
  from_plan_metrics() accepts, with a round-trip test.
- debug_serverless_node_types.py: replace the cell reading the removed
  _captured_plans/_captured_queries internals with the execution/diagnostic API.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AcyKZJKNnbkRiaa44qysbU
Live validation is prepared but unrun: no Databricks credentials were available.
Document the three commands, where the sanitized fixture is printed, and the one
assumption (plan-id propagation into Photon physical nodes) that only a live run
can settle.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AcyKZJKNnbkRiaa44qysbU
Live validation on Databricks serverless failed the `operation ids recorded`
check with None for all six executions. The cause is not Databricks-specific:
ExecutePlanRequest.operation_id is populated only when a caller passes an id
into _execute_plan_request_with_metadata, which the action paths never do, so
the field is an unset optional string and reading it off the request always
yielded empty. The id that identifies the operation is assigned by the server
and arrives on each ExecutePlanResponse.

- Wrap _verify_response_integrity, which the client calls once per response
  before branching on content, so commands carrying no operator metrics record
  an id too. Keep the request-builder hook as a fallback for the caller-supplied
  case. Both stay getattr-guarded, so a client lacking either degrades.
- FakeClient returned a populated op-N id from its request builder, which the
  real client does not; that is why 42 offline cases passed over a hook that
  never fired. The fake now matches the real semantics and emits a FakeResponse.
- Replace the synthetic replay fixture with a real serverless recording (runtime
  4.2.0, Connect client 3.5.0, 40 nodes across two executions, operator names
  sanitized) and assert against its real structure.
- The notebook now reports through whichever channel the outcome uses: notebook
  stdout is not retrievable through the Jobs API, which returns notebook_output
  only on a clean dbutils.notebook.exit and the error field otherwise. Exit with
  the checks, execution summary, and fixture on success; raise with the failed
  checks, their details, and the parse diagnostics on failure.

All 22 notebook checks pass against serverless; 46 offline cases pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The procedure documented as prepared but unrun was run on 2026-09-12 against a
serverless workspace. All 22 checks pass. It settled the two questions offline
tests could not:

- Databricks does not propagate the client-assigned plan id into physical Photon
  nodes, so join keys cannot be matched that way. The documented fallback held:
  both join nodes recorded join_details_source = "unresolved" with a diagnostic
  rather than keys matched by position. Future join attribution work under
  Connect starts from this.
- Operation ids were never recorded on any runtime, fixed in the preceding
  commit.

Also record the remaining degradation, all reported as diagnostics rather than
silent: PhotonRange preserved as an Unknown node, and execute_command actions
producing no operator metrics.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The recording captured two executions that are structurally identical: same
20-node sequence, same 355 metric names, differing only in measured values and
ids. The second exercises no parse path the first does not, so it was 3800 lines
of test data buying a duplicate assertion. Per-query isolation is already covered
against the fake client.

Fixture drops from 7589 to 3799 lines; metric fidelity within the retained
execution is untouched, since pruning the 996 zero-valued entries would make it
no longer a faithful recording.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Setting `_verify_response_integrity = None` on a FakeClient subclass is an
inconsistent override: pyrefly rejects narrowing a method to `Unknown | None`.
Split the stub instead, so the two client shapes are expressed by what each
class has rather than by nulling a method out. LegacyFakeClient predates the
hook; FakeClient adds it. `_emit` resolves the hook through a local, which also
stops it asserting a method the base class does not define.

Caught by CI type check; the earlier local runs covered tests and ruff only.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Both workflows triggered only on pull_request opened/reopened, so any commit
pushed after a PR was opened left it permanently BLOCKED: the required `test`
context had no run on the new head, and workflow_dispatch could not supply one.
A dispatch run's check suite is not associated with the pull request, so its
check lands on the commit but never satisfies the required context -- verified
on PR #20, where `test` and `test-full` both passed on the head sha while the PR
stayed BLOCKED.

Add `merge_group` rather than `synchronize`: the run happens once, when a PR is
queued to merge, instead of on every push, and its check suite is tied to the
merge attempt so it satisfies protection by construction.

Requires the merge queue to be enabled for `main` in Settings → Branches.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Merge queue is not available on this account: the UI offers no "Require merge
queue" option, and creating a ruleset with a `merge_queue` rule is rejected with
"Invalid rule 'merge_queue'" both bare and with full parameters. The merge_group
trigger added in the previous commit would never have fired.

Use `ready_for_review` instead. Open work as a draft and push freely at no CI
cost; marking the PR ready fires one run, and because it comes from a
pull_request event its check suite is associated with the PR and satisfies the
required contexts. `gh pr ready --undo && gh pr ready` re-runs after a later
push.

Also record in the comment that the required contexts are now `test` and
`test-full`, and that the full-object protection PUT drops any setting it does
not restate -- use the granular PATCH on required_status_checks instead.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@jeffbrennan
jeffbrennan marked this pull request as draft September 12, 2026 22:20
@jeffbrennan
jeffbrennan marked this pull request as ready for review September 12, 2026 22:20
@jeffbrennan
jeffbrennan merged commit e6d6847 into main Sep 12, 2026
3 checks passed
@jeffbrennan
jeffbrennan deleted the claude/brief-02-changes-4rzz77 branch September 12, 2026 22: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