Repository navigation
Conversation
A Python operator is not called; it is handed to an interpreter that imports the generated module and drives it through the same open, process and close the engine uses. `PyOpExecHarness` writes that module, starts the driver, and reads back what the operator emitted. The driver is the engine's side of the contract written out plainly: it is what makes the answer this side produces the engine's answer rather than an approximation of it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Automated Reviewer SuggestionsBased on the
|
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #8357 +/- ##
============================================
- Coverage 93.50% 92.55% -0.95%
- Complexity 4882 4964 +82
============================================
Files 1220 1223 +3
Lines 50788 51494 +706
Branches 6262 6327 +65
============================================
+ Hits 47489 47661 +172
- Misses 1746 2242 +496
- Partials 1553 1591 +38
*This pull request uses carry forward flags. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
| config | throughput | MB/s | latency | max Δ latest / 7d | |
|---|---|---|---|---|---|
| 🔴 | bs=10 sw=10 sl=64 | 410 | 0.25 | 23,152/29,109/29,109 us | 🟢 -9.9% / 🔴 +86.9% |
| 🔴 | bs=100 sw=10 sl=64 | 884 | 0.54 | 110,473/161,319/161,319 us | 🔴 +13.5% / 🔴 +52.0% |
| 🔴 | bs=1000 sw=10 sl=64 | 1,085 | 0.662 | 913,909/1,026,857/1,026,857 us | 🔴 +5.6% / 🟢 -5.9% |
Baseline details
Latest main bf356f5 from same runner
| config | metric | PR | latest main | 7d avg | Δ latest | Δ 7d |
|---|---|---|---|---|---|---|
| bs=10 sw=10 sl=64 | throughput | 410 tuples/sec | 441 tuples/sec | 818.88 tuples/sec | -7.0% | -49.9% |
| bs=10 sw=10 sl=64 | MB/s | 0.25 MB/s | 0.269 MB/s | 0.5 MB/s | -7.1% | -50.0% |
| bs=10 sw=10 sl=64 | p50 | 23,152 us | 22,432 us | 12,388 us | +3.2% | +86.9% |
| bs=10 sw=10 sl=64 | p95 | 29,109 us | 32,309 us | 15,825 us | -9.9% | +83.9% |
| bs=10 sw=10 sl=64 | p99 | 29,109 us | 32,309 us | 18,692 us | -9.9% | +55.7% |
| bs=100 sw=10 sl=64 | throughput | 884 tuples/sec | 961 tuples/sec | 1,058 tuples/sec | -8.0% | -16.4% |
| bs=100 sw=10 sl=64 | MB/s | 0.54 MB/s | 0.587 MB/s | 0.646 MB/s | -8.0% | -16.4% |
| bs=100 sw=10 sl=64 | p50 | 110,473 us | 99,649 us | 99,191 us | +10.9% | +11.4% |
| bs=100 sw=10 sl=64 | p95 | 161,319 us | 142,085 us | 106,139 us | +13.5% | +52.0% |
| bs=100 sw=10 sl=64 | p99 | 161,319 us | 142,085 us | 116,669 us | +13.5% | +38.3% |
| bs=1000 sw=10 sl=64 | throughput | 1,085 tuples/sec | 1,106 tuples/sec | 1,087 tuples/sec | -1.9% | -0.2% |
| bs=1000 sw=10 sl=64 | MB/s | 0.662 MB/s | 0.675 MB/s | 0.664 MB/s | -1.9% | -0.3% |
| bs=1000 sw=10 sl=64 | p50 | 913,909 us | 899,943 us | 971,064 us | +1.6% | -5.9% |
| bs=1000 sw=10 sl=64 | p95 | 1,026,857 us | 972,731 us | 1,013,362 us | +5.6% | +1.3% |
| bs=1000 sw=10 sl=64 | p99 | 1,026,857 us | 972,731 us | 1,040,298 us | +5.6% | -1.3% |
Raw CSV
config_idx,batch_size,schema_width,string_len,num_batches,total_ms,total_tuples,total_bytes,tuples_per_sec,mb_per_sec,lat_p50_us,lat_p95_us,lat_p99_us
0,10,10,64,20,487.53,200,128000,410,0.250,23151.56,29109.46,29109.46
1,100,10,64,20,2262.08,2000,1280000,884,0.540,110473.01,161319.08,161319.08
2,1000,10,64,20,18428.82,20000,12800000,1085,0.662,913908.71,1026857.14,1026857.14Eight lines this change is already in the file for. sbt lifts a `.value` written inside a lambda to the top of the task, so the options were read once rather than per suite either way; written where they were, they said the other thing, and sbt warned on it in every build. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Matches apache#8357, which carries this upstream. sbt lifts a `.value` written inside a lambda to the top of the task, so the options were read once either way; written where they were, they said the other thing. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The two helpers here were a second copy of OpExecHarness's, carrying a note that they were kept inline so each harness read on its own and would be consolidated if a third arrived. Two copies of the same twenty lines is already the cost that note was deferring: they have to be changed together, and nothing says so at either site. Preparing the plan does not vary with the executor, so it happens once now. What differs between the harnesses stays here: how the prepared op is run. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Three kinds of comment came out. A drawing of the string the code below assembles. A restatement of a branch the reader can see. And the word MVP, which dated the scope to a moment rather than stating it. What replaces them says the same thing shorter, or says what the code cannot: which cases the harness does not drive and why none of them has an operator asking for it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
carloea2
left a comment
There was a problem hiding this comment.
The Python operator harness looks good.
`PhysicalPlan` is named only by `getPhysicalPlan` and by a comment, so scalafix reports the import and the check fails before any test runs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
carloea2
left a comment
There was a problem hiding this comment.
The unused import was removed. Looks good.
carloea2
left a comment
There was a problem hiding this comment.
Reviewed alongside the related export and verification PRs. These findings are based on code inspection and focused Python checks, not a full Scala suite run.
The driver projected the operator's output onto the schema's attribute names and coerced each value with int/float/str. That is looser than the engine, which hands every yielded value to all_output_to_tuple and finalizes the result: an operator putting 1.5 in an INTEGER field, or a row with a field the schema does not have, or one missing a field, all raise there and none of the three raised here. A reference path that accepts output the engine rejects cannot say what the engine would do, so call the engine's own two functions instead of re-deriving them. finalize pickles a BINARY field's object itself and tags it, where the standalone side pickles the same object untagged, so drop the tag on the way out and both sides carry the same payload. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
carloea2
left a comment
There was a problem hiding this comment.
Further review with local regression tests and connected execution checks.
The driver collected everything the operator yielded and converted the list afterwards. `DataProcessor` reads a yield where it happens, turning it into tuples and finalizing them on the spot, so a UDF that yields one dict twice and changes it in between is recorded as the two values it produced. The driver recorded the last one twice. The operator now takes a callback and hands each value over at the yield, before the generator is advanced. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
DataProcessor writes each input port's schema into the executor before calling on_finish, so a TableOperator handed no rows can still say what its columns were. The driver stands in for DataProcessor here, so it has to do the same or the reference path answers on a frame of no columns where the engine answers on the declared ones. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
It restated its own signature and listed a codec table that lives in TupleIO. The reasons a reader cannot derive stay. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
@carloea2 both findings are fixed: this path runs the engine's own output validation, and each yield is finalized where it is yielded. Would you take another look when you have a moment? |
The handlers are written nine steps later in the set, and this borrowed a fixture writer from them for a single one-column row. Writing it with TupleIO, which the spec already reads with, leaves the order free of the one edge that pointed backwards. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
carloea2
left a comment
There was a problem hiding this comment.
Each emitted tuple is now finalized when it is yielded. This fixes reused mutable output values. Looks good.
carloea2
left a comment
There was a problem hiding this comment.
The harness declares large binary as supported, but both input coercion and output conversion reject it. An identity Python operator with a valid large binary field cannot be checked at all. Please handle large binary and add an input and output case.
The schema sidecar names large_binary and the driver's type map reads it, but neither the input coercion nor the output conversion knew the type, so an operator handed such a column could not be run at all. A large binary is a reference: the bytes live in S3 and the field is the s3:// URI that points at them. Both harnesses now carry that URI, which is what the worker sends and what LargeBinary holds on the other side, and building one touches no storage. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Fixed in f4cdf6a: a large binary is a reference, so both harnesses now carry the s3:// URI the field holds, which is what the worker sends and what LargeBinary holds on the other side. Two cases cover it, one handing the column to the operator as a reference it can read, one writing back a large binary the operator made. |
carloea2
left a comment
There was a problem hiding this comment.
The large binary handling is now covered on both the Python and Scala sides, including input and output. The earlier blocker is fixed. Looks good.
apache#8488 is closed: a table with no rows carries no schema by design, so there is nothing for the driver to hand the operator and no engine behavior here to match. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
carloea2
left a comment
There was a problem hiding this comment.
One input-decoding mismatch reproduced against the actual Arrow tuple provider. Ordinary bytes passed as a control.
The worker's ArrowTableTupleProvider unpickles a BINARY cell that starts with the cast's pickle marker, so a model column reaches the next operator as a model. The driver handed the operator the bytes instead, and a predict on them raised AttributeError. Decoding the cell alone was not enough. The driver then finalized each input tuple, and the cast in finalize pickles any non-bytes value back into bytes. The worker's InputManager attaches the schema without finalizing, so the driver now builds input tuples the same way. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… spec The pickled-cell test used a pickled list in place of a model. It now pickles a fitted DecisionTreeClassifier with the interpreter the harness runs and has the operator call predict on it, which is the call that raised AttributeError when the driver handed over bytes. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…and refuse a row the schema does not describe Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
What changes were proposed in this PR?
A Python operator is not called; it is handed to an interpreter that imports
the generated module and drives it through the same open, process and close
the engine uses.
PyOpExecHarnesswrites that module, starts the driver, andreads back what the operator emitted.
The driver is the engine's side of the contract written out plainly: it is
what makes the answer this side produces the engine's answer rather than an
approximation of it.
A timestamp reaches the operator the way the engine hands one to a Python worker: a naive datetime at microsecond resolution. It is read without pandas, so a year past 2262 arrives intact. That is the zoneless microsecond Arrow type #8704 makes the engine send.
A port that carried no rows reaches the operator with its declared columns, as #8766 makes the engine hand them. The driver declares each input port's schema right before that port finishes, where the worker does.
Eight lines of
build.sbtare unrelated to the export and fixed while this change is in the file: File Service's test grouping read its fork options inside the lambda, which sbt hoists anyway, so the warning it printed on every build said the code meant something it did not.Any related issues, documentation, discussions?
Part of #8325, 4 of 27; that issue lists the set in order. The timestamp handling matches the engine once #8704 lands. It needs #8766, which adds the
input_schemasthe driver writes, so it merges after it.Closes #8409, the task this change is the whole of.
How was this PR tested?
The tests in this change cover it.
PyOpExecHarnessSpecalso checks that a row the output schema does not describe is refused. It also checks that2500-01-01 00:00:00.123456789reaches the operator as2500-01-01 00:00:00.123456. It also checks that a port that carried no rows hands the operator its declared columns. Without the driver's change the operator sees a table with no columns. The whole set is exercised together once the last piece lands: every operator run through the engine and through its generated script, and the two answers compared.Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5, Claude Opus 5.5)