Go SDK: order tasks with Before, After and Label - #74119
Conversation
FrankYang0529
left a comment
There was a problem hiding this comment.
Overall LGTM. Thanks for the PR.
jason810496
left a comment
There was a problem hiding this comment.
Nice, thanks! There're the claude review findings.
jason810496
left a comment
There was a problem hiding this comment.
My one design decision concern before addressing the claude review:
Would it be better to validate the graph at once "after the Dag be registered" so that we don't need to traverse whole graph again whenever a new edge be added?
Just like the TS side that we will only validate the cycle at the "finalize stage". I think we should do the same "last mile validation right before the serve" instead of having iterative validation (since we will compile the SDK into artifact anyway)
A Dag authored in Go can now order tasks that exchange no data, the way
Python's >> and << do, and label the edge that reaches a task.
loaded.Before(notified, cleaned) // loaded >> [notify, cleanup]
cleaned.After(extracted) // cleanup << extracted
loaded.Before(airflow.Label(emptyNotice, "when empty"))
Both verbs are variadic, so one call fans out, and both return their
argument set as one Node, so a.Before(b, c).Before(d) is a >> [b, c] >> d.
airflow.Label wraps the endpoint rather than the call, so each edge of a
fan-out can carry a label of its own. Declaring an edge the Dag already
has only applies the label, which is how an airflow.Inputs edge gets one.
Before and After with no node declare no edge, but they now still reject a Dag that has been registered and a *TaskRef that DagRef.Task did not return, so a mistake in the receiver is not hidden by an empty fan-out. The label a verb puts on an edge is now settled in one place, mergeLabel, which both the check over the call's pairs and the recording step use. An edge verb enumerates its pairs once and records what it checked, so the checks and the recording cannot drift apart, and addEdgeLocked keeps only the cycle check, which cannot run up front. An edge given two labels in one call now says so rather than reporting the first as already on the edge, which no earlier declaration had put there.
Declaring a label on an edge that already carries one now replaces it,
which is what Python's DAG.set_edge_info does: "this will overwrite,
rather than merge with, existing info". It used to panic, a rule the Go
SDK had of its own.
A label on the receiver of an edge verb is still dropped, since Label
marks the edge that reaches a node and the receiver is the node an edge
leaves from. Python is silent there too: Label("x") >> b, with nothing
upstream of the label, sets no label either. Tests now pin both, and the
nesting that Label(Label(x, "inner"), "outer") reads as.
DagRef.Task records the edges Inputs declares so that the cycle check sees them, which is the case ADR-0008 names: b := dag.Task(B, Inputs(a)) then b.Before(a) is a genuine cycle in accepted syntax. Only the upstream and downstream lists were asserted, so deleting that recording failed no test over the cycle itself.
Before, After and Inputs each walked the graph to see whether the edge they were recording closed a cycle, which made building a Dag cost one walk per edge: a 8,000-task chain declared in reverse took seconds. Register now walks the graph once, after the Dag is whole, as the Java SDK's Bundle.register does and the TypeScript SDK plans to. The same chain builds in tens of milliseconds. Recording an edge no longer panics, so a fan-out that is rejected leaves the Dag exactly as it was, and the caveat about a cycle keeping the earlier edges of its call is gone with it. A Dag whose dependencies contain a cycle stays unregistered, so it can still be corrected. A task ordered against itself is still rejected where it is declared: that needs no walk.
653028a to
2effac3
Compare
Agreed, done in the latest push. The cycle check moved into markRegistered, next to the existing Then check, so it runs once over the whole graph instead of once per edge: O(V+E) rather than O(E*(V+E)). I rebuilt your reverse-declared chain locally, 2,000 tasks in 9ms and 8,000 in 24ms end to end (building the Dag plus the Register walk), against the 0.19s/2.1s you measured. Different machine, but it is linear now. This also removes the need for the rollback in your other comment: recording an edge no longer panics, so a rejected fan-out leaves the Dag exactly as it was, and the caveat in the Before doc is gone with it. A task ordered against itself is still rejected where it is declared, since that needs no walk. |
jason810496
left a comment
There was a problem hiding this comment.
Thanks for the update, LGTM now.
After solving the final comment from claude review, it's good to merge.
IfRef.Then and IfRef.Else stored the task on the condition without recording an edge, so the ordering their doc states, that the two run after the condition, was in no Dag the SDK builds: the task had no upstream, and a cycle running through a condition passed Register. read := dag.Task(readRows) loaded := dag.Task(load) dag.If(hasRows, airflow.Inputs(read)).Then(loaded) loaded.Before(read) // read -> hasRows -> load -> read Naming a task now records the edge from the condition to it, which is also the upstream the serialized Dag needs. Declaring that edge again, as an author who writes it out does, changes nothing.
Why
A Dag authored in Go could declare its tasks and the data edges between them
(
airflow.Inputs), but not an ordering between tasks that exchange no data — Python'sa >> bandb << a— and no way to label an edge for the graph view.This implements decisions 6, 7 and 8 of
go-sdk/adr/0008-native-dag-interface.md:airflow.Nodeis the Go counterpart of Python'sDAGNode/DependencyMixin.What
go-sdk/airflow/node.go(new):Node, sealed by an unexportednode(), satisfied by*TaskRef.Before/After, variadic so one call fans out. Both return their argument set as oneNode, not the receiver, soa.Before(b, c).Before(d)isa >> [b, c] >> d.Label(node, text)wraps the endpoint, so each edge of a fan-out can carry its own label.A later label on an edge replaces it, as Python's
DAG.set_edge_infodoes. A label on averb's receiver has no edge to land on and is dropped, as
Label("x") >> bis in Python.airflow.Inputsedge gets one:extracted.Before(Label(transformed, "rows")).go-sdk/airflow/dag.go:DagRef.edgeLabelsholds every edge and its label;TaskRef.upstreams/downstreamshold both sides.DagRef.Taskrecords the edgesInputsdeclares into the same structure.
markRegisteredrejects a Dag whose dependencies contain acycle, next to the existing check for a condition without a
Then: the Dag is whole by then,so one walk of the graph answers for every edge, and the Dag stays unregistered so it can be
corrected. Declaring an edge never walks the graph, so building a Dag is linear in its edges.
Was generative AI tooling used to co-author this PR?