Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces per-cell execution timing for Dataproc Serverless operations in IPython notebooks. It adds an ExecutionTimer class to track and merge round-trip intervals, registers IPython pre- and post-run cell hooks to log timing summaries, and wraps execution methods in ManagedSparkSession to record active durations. Comprehensive unit tests are also added. The feedback suggests improving _timed_iterator to measure only the active waiting time of individual next() calls on the underlying iterator, preventing client-side processing from inflating the reported Dataproc Serverless execution time.
| def _timed_iterator(iterator, timer): | ||
| # This being a generator is load-bearing: `start` is not read until | ||
| # the first next(), so the clock runs from when the caller begins | ||
| # driving the RPC rather than from when the generator was built. | ||
| # Hoisting it out of the generator body would time every fetch as 0s. | ||
| start = time.monotonic() | ||
| try: | ||
| yield from iterator | ||
| finally: | ||
| timer.record(start, time.monotonic()) |
There was a problem hiding this comment.
The current implementation of _timed_iterator measures the entire duration from the first next() call until the generator is fully consumed or closed. However, because this is a generator, any client-side processing (such as iterating over the rows, performing local computations, or plotting) that occurs between pulling chunks from the iterator will be included in the measured duration. This inflates the reported Dataproc Serverless time with client-side execution time, defeating the purpose of separating the two.
To accurately measure only the time spent waiting on the server, we should time each individual next() call on the underlying iterator and accumulate the active waiting time, then record a single interval representing the total active time.
| def _timed_iterator(iterator, timer): | |
| # This being a generator is load-bearing: `start` is not read until | |
| # the first next(), so the clock runs from when the caller begins | |
| # driving the RPC rather than from when the generator was built. | |
| # Hoisting it out of the generator body would time every fetch as 0s. | |
| start = time.monotonic() | |
| try: | |
| yield from iterator | |
| finally: | |
| timer.record(start, time.monotonic()) | |
| def _timed_iterator(iterator, timer): | |
| # This being a generator is load-bearing: `start` is not read until | |
| # the first next(), so the clock runs from when the caller begins | |
| # driving the RPC rather than from when the generator was built. | |
| # Hoisting it out of the generator body would time every fetch as 0s. | |
| active_time = 0.0 | |
| try: | |
| while True: | |
| start = time.monotonic() | |
| try: | |
| val = next(iterator) | |
| except StopIteration: | |
| break | |
| finally: | |
| active_time += time.monotonic() - start | |
| yield val | |
| finally: | |
| end = time.monotonic() | |
| timer.record(end - active_time, end) |
445d814 to
7565e39
Compare
Adds the timing primitive: module level state accumulating what a cell spends in Managed Spark, reported against the cell's own wall time. Each round trip is kept with its kind and duration, alongside session creation and the bytes pulled over the websocket bridge, so the line distinguishes a slow query from a cold session, a schema lookup, or a large result dragging back through the tunnel. Shared state rather than locals, because the timestamps are taken in different call stacks. Round trips start and end inside gRPC wrappers on threads that are not necessarily the main one, transport is counted on the bridge's forwarding threads, and the reporting happens in an IPython cell callback. Nothing sees all of it, so the totals sit behind a lock. Analyze is a kind of round trip rather than a concept of its own, which is why record takes a kind and there is no separate recorder for it. Anything recorded while no cell is in progress is dropped, keeping this inert and bounded outside IPython, where nothing ever calls start_cell. Round trips are clipped to the cell window, so a generator built in an earlier cell and drained in this one cannot report time from before the cell began. Both the Managed Spark total and the transport blocked time are clamped to the cell's own duration: concurrent round trips and concurrent forwarding threads each accumulate into one counter, so either could otherwise sum past the cell that contains it. register_cell_timing starts a cell when none is in progress. On the first cell nothing has registered a pre_run_cell handler yet, so without this the cell that creates the session, and pays the largest cost there is, would report nothing at all.
Feeds the timer from the Spark Connect client and the websocket bridge,
and registers the IPython hooks, so a cell that touches the cluster
logs what it cost:
Cell took 61.79s, 61.79s in Managed Spark
(session creation 61.79s)
Cell took 1.67s, 1.67s in Managed Spark across 1 round trip
(analyze 1.67s)
Cell took 23.46s, 22.48s in Managed Spark across 1 round trip
(fetch 22.48s; transport 62.1 MiB down in 23.32s, 2.7 MiB/s)
Three call sites do the work. _analyze was previously untimed, so a
cell that only read a schema reported nothing despite costing 1.67s
against a live cluster. Session creation was already measured but
logged at debug, which basicConfig(INFO) suppresses, hiding the single
largest cost in a first cell; it is now recorded against that cell and
logged at info so batch runs see it too. The bridge's recv and send
count payload bytes and blocked time, which is how the 2.7 MiB/s above
became visible at all.
register_cell_timing is also called at the top of getOrCreate, before
any creation work, so provisioning lands inside a live cell rather than
being dropped. It is idempotent, so the existing call on the reuse path
stays as it is.
_execute is blocking, so try/finally around it measures the round trip.
_execute_and_fetch_as_iterator is a generator function, so the wrapper
returns a wrapping generator; timing the call itself would record zero
seconds for every collect. All paths record in a finally, because a
round trip that failed or was abandoned still spent the time.
Lazy SELECT plan registrations are round trips too, so they are timed
and counted. Only their link display stays suppressed, as before.
7565e39 to
1cdaae6
Compare
What
Cells that touch the cluster log what they cost. These are real lines from real IPython cells against a live session, not examples:
A cell doing pure client-side work logs nothing.
Each round trip is named by kind (
execute,fetch,analyze) and timed individually — the three slowest are listed, with+N morebeyond that. Note line two: a singlecollect()is two fetch round trips, which an aggregate count would have hidden.Why
A throwaway spike measured where the time actually goes before any of this was designed:
.schema_analyzewas untimedTransport was ~95% of the only scenario that hurt, and the bridge hex-encodes (
proxy.py), so 62 MiB of Arrow crossed the wire as 124 MiB of hex text. Throughput held at 2.5-3.0 MiB/s across four independent sessions. Without these counters a report like that just looks like "Spark is slow".Jupyter, Colab and
%%timealready give you the first number in each line. None of them can give you the rest.How
execution_timer.py— module-level state, no classes. A list of(kind, seconds)round trips plus three counters, all behind one lock, fed from three places and read by the IPython cell hooks.start_cell(), so the state stays bounded for the life of a batch job.register_cell_timing()is idempotent; it's called per session creation and IPython'sevents.registeraccepts duplicates silently.session.pywraps_analyze,_executeand_execute_and_fetch_as_iterator, and feeds session creation.proxy.pycounts payload bytes and blocked time inbridged_socket.recv/send.Three bugs found by running it for real
The first cell reported nothing — the one paying the 56s. IPython fires
pre_run_cellfor cell 1 before anything has registered a handler, so no cell was marked started andrecord_session_creationhit the drop-when-no-cell rule. Every unit test passed while the feature was inert in its most important case. Fixed by registering at the top ofgetOrCreate()and having registration start a cell when none is in progress. Caveat: cell 1'scell_secondsis then measured from registration rather than true cell start, so it slightly under-reports its own wall time.transport 0.0 MiB down in 43.24s, 0.0 MiB/s— sub-MiB control traffic printed as two zeros and read as broken. Now reports KiB and drops the meaningless rate, keeping the real signal that 43s went by blocked while moving almost nothing (a cold cluster).0.00s in Dataproc Serverless across 0 operations— the original first-cell line was self-contradictory: it claimed zero time on a service where it had just spent a minute. Session creation now counts toward the headline, "operations" became "round trips", and the count is omitted rather than printed as zero.The subtle bit
_execute_and_fetch_as_iteratoris a generator function — calling it returns a generator without sending anything, so timing the call records ~0s for everydf.collect(). The wrapper returns a wrapping generator, and because_timed_iteratoris itself a generator,startisn't read until the firstnext(). There's a comment saying the generator-ness is load-bearing, and a regression test asserting zero round trips are recorded when the wrapper is called but not consumed.time.monotonic()throughout, deviating fromtime.time()elsewhere insession.py— those are wall-clock deadlines, these are durations.Decisions worth a second opinion
record()takes a kind; there's norecord_analyze(). A schema-only cell must count as a round trip, or it stays silent and the blind spot survives the fix.logger.debugfor session creation was promoted tologger.info.basicConfig(level=logging.INFO)was suppressing a 56-second event, and outside IPython the cell line never fires, so info is the only way a batch user learns of it.num_bytes_readwas considered and dropped. The spike read 0.0 MiB on every stage — but every scenario usedspark.range(), which reads no files, so that's untested rather than disproven.Testing
pytest tests/unit→ 227 passed, 13 subtests passed, 1 failed. The failure is pre-existing and environmental:test_create_session_without_application_default_credentialsfails on any machine where ADC are configured.pyink --checkclean.Plus four live runs against real sessions, which is how all three bugs above were found.