You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Proposal: operators that emit on more than one output port
#8914
A common step in an analysis produces more than one thing: a cleaned table and a chart of it, a train split and a test split, matched rows and rejects. Today one operator shows one result. To get a table and a plot you chain a second operator off the same link and look at two result tabs. That works, but it adds another operator for each additional output and a Python UDF that already has both in hand has no way to hand them out separately.
The port model already supports this. Split and If emit on two ports today. What is missing is the plumbing between a Python UDF and the result panel. This proposal is to finish that plumbing, in small PRs, and to agree on the API decisions up front.
Current architecture
Ports are first-class from the descriptor down to storage:
flowchart LR
D[OperatorInfo<br/>outputPorts: List OutputPort<br/>id, displayName, mode] --> P[PhysicalOp<br/>outputPorts: Map port -> schema, links]
P --> C[WorkflowCompiler<br/>storage for every external port<br/>of a viewed op]
C --> S[Scheduler<br/>AssignPort per output port<br/>with that port's schema]
S --> W[Worker]
W --> ST[(Storage<br/>URI keyed by op + port)]
Loading
The Scala worker routes by port. An executor returns (tuple, Some(port)) from processTupleMultiPort, and OutputManager.passTupleToDownstream sends it only on links whose fromPortId matches. Split and If use this.
Two paths stop honoring the port.
Python worker. The UDF yields tuples with no port. pyamber finalizes each tuple against the first port's schema, then hands it to every partitioner, so it goes out on every link of every port and into every port's storage.
flowchart LR
U[UDF yields tuple] --> F[finalize with<br/>port 0 schema]
F --> O[OutputManager.tuple_to_batch<br/>all partitioners]
O --> L0[links off port 0]
O --> L1[links off port 1]
F --> SS[storage of every port]
Loading
Result delivery. Result delivery, pagination, export, and sync execution then select port 0. The websocket result event is keyed by operator, and the result panel builds one Result tab per operator. The output mode (table vs HTML) is taken from the first external port encountered, so one operator is either all table or all HTML.
flowchart LR
ST[(storage<br/>op + port 0)] --> R[ExecutionResultService<br/>getResultUri op, PortIdentity 0]
ST1[(storage<br/>op + port 1)] -. not delivered independently .-> R
R --> E[WebResultUpdateEvent<br/>Map opId -> update]
E --> UI[result panel<br/>one Result tab per operator]
Loading
Summary of what exists and what does not:
Layer
Per-port today
Port model, schema propagation, compiler, scheduler, storage
yes
Scala executor API and output manager
yes (Split, If)
Frontend: add/remove ports, persistence, links by port, per-port row counts on the canvas
yes
Python UDF descriptor: ports declared by the user
yes
Python UDF descriptor: output schema for ports after the first
no, port 0 only
pytexera: a way for code to name a port
no
pyamber: routing, schema, stats, storage by port
no, port 0 / all
Result service, export, sync API, websocket event, result panel
no, port 0 only
Related: #8274 and #8276 (pyamber get_port ignores the id; fixes the lookup only, nothing calls it with a port yet), #2078 (UDF port validation).
What would change
Four incremental PRs: the first delivers multi-port results for Scala operators; the remaining PRs build toward multi-port Python UDF output.
1. Results per port
Key the result store by operator plus port. Send the websocket update per port, add a portId (default 0) to the pagination, export and sync requests, and read the output mode from the port being served. The result panel shows one tab per external output port when an operator has more than one, labeled by the port's display name.
flowchart LR
ST0[(op + port 0)] --> R[ExecutionResultService<br/>keyed by op + port]
ST1[(op + port 1)] --> R
R --> E[WebResultUpdateEvent<br/>per op + port]
E --> UI[result panel<br/>Result: table / Result: chart]
Loading
Single port workflows retain their behavior and presentation, so existing workflows should see no difference. This alone makes both ports of Split and If viewable.
2. pyamber honors a port
tuple_to_batch(tuple, port) filters partitioners by the link's source port. data_processor finalizes against that port's schema. Statistics and storage writes take the port. Lands on top of #8276. Arrow batch construction uses that port’s schema, including buffered batches flushed at completion.
3. pytexera: naming the port
Add an explicit PortOutput wrapper that carries a value and its destination port. Bare yields continue to target port 0. The SDK will introduce a return type covering both bare and wrapped outputs.
A bare yield targets port 0. all_output_to_tuple learns to unwrap PortOutput, normalize its value into tuples, and preserve the destination port for each tuple.
4. Per-port output schema and mode on the Python UDF
The descriptor propagates a schema for port 0 only, and the saved port description carries no schema or mode. Add a per-port kind:
kind
schema
mode
table (default)
retain input columns + extra output columns, as today
SET_SNAPSHOT
visualization
html-content: STRING
SINGLE_SNAPSHOT
Workflows saved before this default every port to table. The port editor in the property panel gets the kind field for output ports; today it edits input ports only.
Nothing changes in the scheduler, storage layout, Arrow data path (maybe Arrow serialization code?), or the PythonDataHeader proto: partitioning happens on the Python side, so Scala never needs the port.
Where I want feedback
Unnamed output on a multi-port Python UDF. Scala treats "no port" as "all ports" for sending and asserts a single port for schema. I propose Python defaults to port 0, so a chart row is never validated against the table schema. Should Scala be aligned, or is "all ports" a semantic anyone relies on?
** Explicit PortOutput wrapper vs. pair vs. emit.** I propose yield PortOutput(output, port=1) because it preserves the generator contract and makes the destination explicit. Would a raw (output, port) pair or self.emit(output, port=1) be preferable?
Per-port kind vs. per-port column list.kind covers the table+chart case with one field. A per-port extra-columns list is more general but a bigger form change. Start with kind and add columns later?
Websocket shape. Nested Map[opId, Map[portId, update]] vs. a composite key "opId:portId". The composite key is the smaller frontend diff; the nested map keeps one shape for every operator and avoids key parsing in every consumer
Acceptance Criteria
Validate a UDF with at least three outputs, including different table/HTML schemas, an empty output port, and buffered output flushed at completion. Verify independent results, pagination, export, invalid-port errors, and compatibility with existing single-port workflows.
Note - This supports any number of output ports; table plus visualization is the initial use case. Independently configurable table schemas are future work.
reacted with thumbs up emoji reacted with thumbs down emoji reacted with laugh emoji reacted with hooray emoji reacted with confused emoji reacted with heart emoji reacted with rocket emoji reacted with eyes emoji
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Why
A common step in an analysis produces more than one thing: a cleaned table and a chart of it, a train split and a test split, matched rows and rejects. Today one operator shows one result. To get a table and a plot you chain a second operator off the same link and look at two result tabs. That works, but it adds another operator for each additional output and a Python UDF that already has both in hand has no way to hand them out separately.
The port model already supports this. Split and If emit on two ports today. What is missing is the plumbing between a Python UDF and the result panel. This proposal is to finish that plumbing, in small PRs, and to agree on the API decisions up front.
Current architecture
Ports are first-class from the descriptor down to storage:
The Scala worker routes by port. An executor returns
(tuple, Some(port))fromprocessTupleMultiPort, andOutputManager.passTupleToDownstreamsends it only on links whosefromPortIdmatches. Split and If use this.Two paths stop honoring the port.
Python worker. The UDF yields tuples with no port. pyamber finalizes each tuple against the first port's schema, then hands it to every partitioner, so it goes out on every link of every port and into every port's storage.
Result delivery. Result delivery, pagination, export, and sync execution then select port 0. The websocket result event is keyed by operator, and the result panel builds one Result tab per operator. The output mode (table vs HTML) is taken from the first external port encountered, so one operator is either all table or all HTML.
Summary of what exists and what does not:
Related: #8274 and #8276 (pyamber
get_portignores the id; fixes the lookup only, nothing calls it with a port yet), #2078 (UDF port validation).What would change
Four incremental PRs: the first delivers multi-port results for Scala operators; the remaining PRs build toward multi-port Python UDF output.
1. Results per port
Key the result store by operator plus port. Send the websocket update per port, add a
portId(default 0) to the pagination, export and sync requests, and read the output mode from the port being served. The result panel shows one tab per external output port when an operator has more than one, labeled by the port's display name.Single port workflows retain their behavior and presentation, so existing workflows should see no difference. This alone makes both ports of Split and If viewable.
2. pyamber honors a port
tuple_to_batch(tuple, port)filters partitioners by the link's source port.data_processorfinalizes against that port's schema. Statistics and storage writes take the port. Lands on top of #8276. Arrow batch construction uses that port’s schema, including buffered batches flushed at completion.3. pytexera: naming the port
Add an explicit PortOutput wrapper that carries a value and its destination port. Bare yields continue to target port 0. The SDK will introduce a return type covering both bare and wrapped outputs.
A bare yield targets port 0. all_output_to_tuple learns to unwrap PortOutput, normalize its value into tuples, and preserve the destination port for each tuple.
4. Per-port output schema and mode on the Python UDF
The descriptor propagates a schema for port 0 only, and the saved port description carries no schema or mode. Add a per-port
kind:table(default)visualizationhtml-content: STRINGWorkflows saved before this default every port to
table. The port editor in the property panel gets the kind field for output ports; today it edits input ports only.Nothing changes in the scheduler, storage layout, Arrow data path (maybe Arrow serialization code?), or the
PythonDataHeaderproto: partitioning happens on the Python side, so Scala never needs the port.Where I want feedback
kindvs. per-port column list.kindcovers the table+chart case with one field. A per-port extra-columns list is more general but a bigger form change. Start withkindand add columns later?Map[opId, Map[portId, update]]vs. a composite key"opId:portId". The composite key is the smaller frontend diff; the nested map keeps one shape for every operator and avoids key parsing in every consumerAcceptance Criteria
Validate a UDF with at least three outputs, including different table/HTML schemas, an empty output port, and buffered output flushed at completion. Verify independent results, pagination, export, invalid-port errors, and compatibility with existing single-port workflows.
Note - This supports any number of output ports; table plus visualization is the initial use case. Independently configurable table schemas are future work.
All reactions