uranusjr commented on code in PR #74230:
URL: https://github.com/apache/airflow/pull/74230#discussion_r4203701996
##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt:
##########
@@ -133,12 +143,122 @@ class DagDef(
task.owner?.let { owner ->
throw IllegalArgumentException("Task '${task.id}' already belongs to Dag
'${owner.id}'")
}
+ require(task.id !in groups) { "Dag '$id' already has a task group with ID:
${task.id}" }
require(tasks.putIfAbsent(task.id, task) == null) {
"Tasks in Dag have duplicate ID: ${task.id}"
}
task.owner = this
return this
}
+
+ /**
+ * Declares a task group of this Dag.
+ *
+ * ```java
+ * var staging = dag.taskGroup("staging");
+ * var stage = staging.task("stage", Stage.class); // task "staging.stage"
+ * extract.before(staging);
+ * ```
+ *
+ * @param id Group ID. Must contain only ASCII letters, digits, underscores,
+ * or dashes, and differ from every task and group ID in this Dag.
+ * @return The group, to declare tasks in and to wire edges with.
+ * @throws IllegalArgumentException if [id] is not a valid group ID, or the
+ * Dag already has a task or task group with that ID.
+ */
+ fun taskGroup(id: String): TaskGroupRef = addGroup(null, id)
+
+ internal fun addGroup(
+ parent: TaskGroupRef?,
+ localId: String,
+ ): TaskGroupRef {
+ require(GROUP_ID.matches(localId)) {
+ "Task group ID '$localId' must contain only ASCII letters, digits,
underscores, or dashes"
+ }
+ val groupId = parent?.qualify(localId) ?: localId
+ require(groupId !in tasks && groupId !in groups) {
+ "Dag '$id' already has a task or task group with ID: $groupId"
+ }
+ return TaskGroupRef(this, groupId, parent).also {
+ groups[groupId] = it
+ parent?.children?.add(it)
+ }
+ }
+
+ /**
+ * Turns every edge drawn to or from a task group into edges between tasks,
+ * and records it on the group as Python's `TaskGroup` does.
+ *
+ * A group upstream stands for its leaves and a group downstream for its
+ * roots. Edges expand in the order they were drawn, each reading the task
+ * edges the ones before it left behind, which is how Python resolves a
+ * group's endpoints at every `>>`. Runs once: a group edge drawn after
Review Comment:
Actually…
Python resolves a group's endpoints when `>>` is called, using only the
tasks the group holds at that moment. This PR resolves them when the Dag is
registered, using the group's final contents, so a task added to a group after
an edge is drawn still gets wired. We should either fix the implementation or
document the difference.
--
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]