uranusjr commented on code in PR #74230:
URL: https://github.com/apache/airflow/pull/74230#discussion_r4203677237
##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt:
##########
@@ -98,9 +98,87 @@ interface Deps {
fun <T> lit(value: T?): Arg<T> = Arg.lit(value)
}
-/** Several tasks as one point in the flow, which no single [TaskRef] can
represent. */
+/**
+ * One task group of the Dag being wired: a point in the flow, and the
+ * namespace of the tasks and groups declared inside it.
+ *
+ * The generated wiring view nests one of these per [Builder.TaskGroup] class,
+ * so a group is reached by calling it and its contents by calling on through:
+ *
+ * ```java
+ * staging().stage(rows); // the task "staging.stage"
+ * staging().checks().nulls(id); // the task "staging.checks.nulls"
+ * extract().before(staging()); // the whole group runs after extract
+ * ```
+ */
+interface Group : Deps.Flow {
+ /** Full ID of this group, as the Dag registered it. */
+ fun groupId(): String
+
+ override fun nodes(): List<TaskDef> = Refs.group(groupId()).nodes()
+}
+
+/** Several tasks or groups as one point in the flow, which no single handle
can represent. */
internal class FlowSet(
- private val nodes: List<TaskDef>,
+ internal val flows: List<Deps.Flow>,
) : Deps.Flow {
- override fun nodes(): List<TaskDef> = nodes
+ override fun nodes(): List<TaskDef> = flows.flatMap { it.nodes() }
+}
+
+/**
+ * Draws an ordering edge from each endpoint of [upstream] to each of
+ * [downstream]. An edge between two tasks is recorded on the downstream task;
+ * one with a task group at either end is recorded on the group's Dag and
+ * expanded onto tasks when the Dag is registered, once the group is complete.
+ */
+private fun link(
+ upstream: Deps.Flow,
+ downstream: Deps.Flow,
+) {
+ for (up in upstream.endpoints()) {
+ for (down in downstream.endpoints()) {
+ if (up is TaskDef && down is TaskDef) {
Review Comment:
Instead of passing in Any and use unchecked casts, we can probably add a
real type instead, e.g. `internal sealed Endpoint` that TaskDef and
TaskGroupRef implement. Each Flow can provide its endpoints polymorphically,
not through a when here. This also lets `DagDef.groupEdges` stop being
`Pair<Any, Any>`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]