Skip to content

test(verify): run a Python operator the way the engine runs it - #8357

Draft
kz930 wants to merge 17 commits into
apache:mainfrom
kz930:feat/verify-run-operator-python
Draft

kz930 wants to merge 17 commits into
apache:mainfrom
kz930:feat/verify-run-operator-python

Conversation

@kz930

@kz930 kz930 commented Sep 2, 2026 •

Copy link
Copy Markdown
Contributor

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. 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.

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.sbt are 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_schemas the 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. PyOpExecHarnessSpec also checks that a row the output schema does not describe is refused. It also checks that 2500-01-01 00:00:00.123456789 reaches the operator as 2500-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)

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>
@github-actions github-actions Bot added feature dependencies Pull requests that update a dependency file common platform Non-amber Scala service paths labels Sep 2, 2026
@github-actions

github-actions Bot commented Sep 2, 2026 •

Copy link
Copy Markdown
Contributor

Automated Reviewer Suggestions

Based on the git blame history of the changed files, we recommend the following reviewers:

  • Contributors with relevant context: @aicam
    You can notify them by mentioning @aicam in a comment.

@codecov-commenter

codecov-commenter commented Sep 2, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 92.55%. Comparing base (bf356f5) to head (135b9be).
⚠️ Report is 90 commits behind head on main.

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     
Flag Coverage Δ *Carryforward flag
access-control-service 77.38% <ø> (+5.59%) ⬆️
agent-service 99.32% <ø> (ø) Carriedforward from a2f9668
amber 88.06% <ø> (-1.04%) ⬇️
computing-unit-managing-service 60.41% <ø> (-16.74%) ⬇️
config-service 87.37% <ø> (+0.12%) ⬆️
file-service 81.53% <ø> (ø)
frontend 96.72% <ø> (+0.02%) ⬆️ Carriedforward from a2f9668
notebook-migration-service 83.73% <ø> (ø)
pyamber 98.47% <ø> (ø) Carriedforward from a2f9668
workflow-compiling-service 74.09% <ø> (ø)

*This pull request uses carry forward flags. Click here to find out more.

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@github-actions

github-actions Bot commented Sep 2, 2026 •

Copy link
Copy Markdown
Contributor

⚠️ Benchmark changes need a look

🟢 2 better · 🔴 9 worse · ⚪ 4 noise (<±5%) · 0 without baseline

Compared against main bf356f5 benchmarked on this same runner, so the delta is largely free of cross-runner hardware noise. The "7d avg" column still reflects the gh-pages dashboard. Treat <±5% as noise unless repeated.

Dashboard · Run

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.14

Eight 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>
kz930 added a commit to Nicoleee1108/texera_workflow_to_py that referenced this pull request Sep 3, 2026
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>
kz930 and others added 2 commits September 3, 2026 15:11
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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The Python operator harness looks good.

@kz930
kz930 marked this pull request as draft September 4, 2026 17:44
`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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The unused import was removed. Looks good.

@carloea2 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Further review with local regression tests and connected execution checks.

Comment thread workflow-compiling-service/src/test/resources/python/py_op_driver.py Outdated
kz930 and others added 2 commits September 10, 2026 13:02
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>
@kz930

kz930 commented Sep 11, 2026 •

Copy link
Copy Markdown
Contributor Author

@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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Each emitted tuple is now finalized when it is yielded. This fixes reused mutable output values. Looks good.

@carloea2 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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>
@kz930

kz930 commented Sep 18, 2026

Copy link
Copy Markdown
Contributor Author

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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 carloea2 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One input-decoding mismatch reproduced against the actual Arrow tuple provider. Ordinary bytes passed as a control.

Comment thread workflow-compiling-service/src/test/resources/python/py_op_driver.py Outdated
kz930 and others added 2 commits September 23, 2026 16:16
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>
@kz930

kz930 commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

@carloea2 A Python operator now receives a timestamp as the engine sends it after #8704: a naive datetime at microsecond resolution.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

common dependencies Pull requests that update a dependency file feature platform Non-amber Scala service paths

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Run a Python operator the way the engine runs it

3 participants