jason810496 commented on code in PR #74230:
URL: https://github.com/apache/airflow/pull/74230#discussion_r4204465451


##########
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 {

Review Comment:
   Moved it under `Deps` as `Deps.TaskGroup`, beside `Flow`, so nothing with a 
generic name lands at the top level of `org.apache.airflow.sdk`. The generated 
view inherits the nested name, so a group view reads `interface Staging extends 
Deps.TaskGroup`.
   
   Fixed in a22621f5749.



##########
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:
   Done. Added `sealed interface Endpoint`, implemented by `TaskDef` and 
`TaskGroupRef`. `endpoints(): List<Endpoint>` is now a member of `Flow` 
defaulting to `nodes()`, overridden by `TaskGroupRef`, `Deps.TaskGroup` and 
`FlowSet`, so the `when` in `link` is gone. `DagDef.groupEdges` is 
`Pair<Endpoint, Endpoint>` and no unchecked cast is left.
   
   Fixed in a22621f5749.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt:
##########
@@ -241,3 +361,5 @@ interface Task {
     client: Client,
   )
 }
+
+private val GROUP_ID = Regex("[A-Za-z0-9_-]+")

Review Comment:
   Both now use one copy, in `org.apache.airflow.sdk.internal`, which is where 
the processor already imports `SchemaFields` from.
   
   Fixed in a22621f5749.



##########
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:
   I intentionally make all the SDKs so far to resolve them all once.
   
   Noted in the task-groups section of `java.rst` and in `TaskGroupRef`'s KDoc: 
endpoints are read once, from everything the group holds by then.
   
   Documented in a22621f5749.



##########
java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt:
##########
@@ -682,13 +829,45 @@ class BuilderProcessor : AbstractProcessor() {
   }
 }
 
-/** One [Builder.Task]-annotated method with its resolved id and data 
parameters. */
+/** The tasks and task groups one class declares. */
+private class Scope(
+  val tasks: List<TaskDeclaration>,
+  val groups: List<GroupDeclaration>,
+) {
+  /** Every task of this scope and the groups beneath it, outermost first. */
+  fun allTasks(): List<TaskDeclaration> = tasks + groups.flatMap { 
it.scope.allTasks() }
+
+  /** Every group beneath this scope, parents before the groups nested in 
them. */
+  fun allGroups(): List<GroupDeclaration> = groups.flatMap { listOf(it) + 
it.scope.allGroups() }
+}
+
+/** One `@Builder.TaskGroup` class, and what it declares. */
+private class GroupDeclaration(
+  val element: TypeElement,
+  val id: String,
+  val fullId: String,
+  val scope: Scope,
+) {
+  /** The view method that reaches this group, named after the class it is 
declared as. */
+  val accessor: String = 
element.simpleName.toString().replaceFirstChar(Char::lowercase)
+}
+
+/**
+ * One [Builder.Task]-annotated method with its resolved id and data 
parameters.
+ * [owner] is the class that declares it, which the generated body 
instantiates,
+ * and [classPath] the task-group classes enclosing it.
+ */
 private class TaskDeclaration(
   val method: ExecutableElement,
   val id: String,
   val dataParams: List<DataParam>,
+  val owner: TypeElement,
+  val classPath: List<String> = emptyList(),
+  /** Full ID of the group this task sits in, empty when it sits in none. */
+  val groupId: String = "",
 ) {
-  val className: String = 
method.simpleName.toString().replaceFirstChar(Char::uppercase)
+  val className: String =
+    (classPath + 
method.simpleName.toString().replaceFirstChar(Char::uppercase)).joinToString("_")

Review Comment:
   Detecting and reporting for now: `A_B` holding task `c`, and `A` holding `B` 
holding `c`, both generate `A_B_C` and now fail the build with both method 
names in the message.
   
   Detection added in a22621f5749.



##########
java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt:
##########
@@ -717,13 +896,19 @@ private val REFS_TYPE = ClassName.get(Refs::class.java)
 private val ARG_TYPE = ClassName.get(Arg::class.java)
 private val TASK_HANDLE_TYPE = ClassName.get(TaskRef::class.java)
 private val DEPS_TYPE = ClassName.get(Deps::class.java)
+private val GROUP_TYPE = ClassName.get(Group::class.java)
+private val LIST_TYPE = ClassName.get(List::class.java)
+private val MAP_TYPE = ClassName.get(Map::class.java)
+
+private val GROUP_ID = Regex("[A-Za-z0-9_-]+")
 
 private const val DAG_ANNOTATION = "org.apache.airflow.sdk.Builder.Dag"
 private const val TASK_ANNOTATION = "org.apache.airflow.sdk.Builder.Task"
 
 private val RESERVED_VIEW_NAMES =
   setOf(
     "depends",
+    "group",

Review Comment:
   Removed. Nothing on `Deps` or on the generated view is named `group`, so it 
only rejected valid code.
   
   Fixed in a22621f5749.



##########
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
+   * that is rejected where it is drawn rather than silently left out.
+   */
+  internal fun expandGroupEdges() {
+    if (groupEdgesExpanded) return
+    groupEdgesExpanded = true
+    for ((upstream, downstream) in groupEdges) {
+      val from = leavesOf(upstream)
+      rootsOf(downstream).forEach { it.dependsOn(*from.toTypedArray()) }
+      if (downstream is TaskGroupRef) {
+        downstream.upstreamTaskIds += from.map { it.id }
+        if (upstream is TaskGroupRef) downstream.upstreamGroupIds += 
upstream.id
+      }
+      // When both ends are groups, the upstream records the downstream group
+      // only, not its tasks, which is how Python leaves it.
+      if (upstream is TaskGroupRef) {
+        if (downstream is TaskGroupRef) {
+          upstream.downstreamGroupIds += downstream.id
+        } else {
+          upstream.downstreamTaskIds += (downstream as TaskDef).id
+        }
+      }
+    }
+  }
+
+  /**
+   * The tasks an edge out of [endpoint] starts from: Python's `find_leaves`.
+   *
+   * A group stands for its own leaves. One holding none stands for whatever
+   * already runs before it, then for the group edges drawn into it, and
+   * failing both for the group it is nested in, so an edge out of an empty
+   * group still reaches the tasks around it.
+   */
+  private fun leavesOf(
+    endpoint: Any,
+    seen: MutableSet<TaskGroupRef> = mutableSetOf(),
+  ): List<TaskDef> {
+    if (endpoint is TaskDef) return listOf(endpoint)
+    var group: TaskGroupRef? = endpoint as TaskGroupRef
+    while (group != null && seen.add(group)) {
+      group.leaves().takeIf { it.isNotEmpty() }?.let { return it }
+      group.upstreamTaskIds.takeIf { it.isNotEmpty() }?.let { ids -> return 
ids.map { tasks.getValue(it) } }
+      val current = group
+      groupEdges.filter { it.second === current }.takeIf { it.isNotEmpty() 
}?.let { drawn ->
+        return drawn.flatMap { leavesOf(it.first, seen) }
+      }

Review Comment:
   Good catch.`leavesOf` is now leaves, then `upstreamTaskIds`, then the 
parent, and `rootsOf` is just `roots()`.
   
   Fixed in a22621f5749.



##########
java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt:
##########
@@ -326,22 +389,75 @@ class BuilderProcessor : AbstractProcessor() {
   private fun inType(paramType: TypeMirror): TypeName =
     ParameterizedTypeName.get(ARG_TYPE, 
WildcardTypeName.subtypeOf(TypeName.get(paramType).boxIfPossible()))
 
-  private fun collectTasks(el: TypeElement): List<TaskDeclaration> {
-    val declarations = mutableListOf<TaskDeclaration>()
+  /** The Dag's tasks and task groups, read from the class tree the author 
wrote. */
+  private fun collectScope(
+    el: TypeElement,
+    path: List<String>,
+    classPath: List<String>,
+  ): Scope {
+    val tasks = mutableListOf<TaskDeclaration>()
     for (inner in el.enclosedElements) {
       if (inner !is ExecutableElement) continue
       val ann = inner.getAnnotation(Builder.Task::class.java) ?: continue
       if (inner.isVarArgs) throw IllegalArgumentException("Cannot create task 
from vararg function ${inner.simpleName}")
-      val id = ann.id.ifBlank { inner.simpleName.toString() }
-      require(declarations.none { it.id == id }) { "Tasks in Dag have 
duplicate ID: $id" }
-      require(declarations.none { 
it.method.simpleName.contentEquals(inner.simpleName) }) {
-        "Dag class ${el.simpleName} overloads task method 
'${inner.simpleName}'; a method's name is the " +
+      val localId = ann.id.ifBlank { inner.simpleName.toString() }
+      require('.' !in localId) {
+        "Task ID '$localId' on method '${inner.simpleName}' of 
${el.simpleName} contains '.', which " +
+          "Airflow reads as a task group prefix; declare the task inside a 
@Builder.TaskGroup class instead"
+      }

Review Comment:
   Removed. Python's `KEY_REGEX` is `^[\w.-]+$`, and group membership now comes 
from the class tree rather than the text of the id, so the check only rejected 
valid code. Added a test that `@Builder.Task(id = "staging.stage")` compiles.
   
   Fixed in a22621f5749.



##########
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
+   * that is rejected where it is drawn rather than silently left out.
+   */
+  internal fun expandGroupEdges() {

Review Comment:
   Done. `expandGroupEdges()` now returns a `GroupExpansion` and stores 
nothing: `TaskGroupRef` lost its four mutable id sets and `DagDef` lost the 
expanded flag. `Bundle.register`'s cycle check computes it on demand, and 
#71190's serializer does the same.
   
   Fixed in a22621f5749, with the serializer side in #71190.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt:
##########
@@ -57,9 +59,30 @@ object Refs {
     dag: DagDef,
     taskIds: List<String>,
     depends: Runnable,
+  ): DagDef = record(dag, taskIds, emptyList(), emptyMap(), depends)
+
+  /**
+   * Runs one `depends()` call as [record] does, for a Dag whose tasks sit in
+   * task groups.
+   *
+   * @param groupIds Full ID of every task group, parents before the groups
+   *    nested in them. All are created before `depends()` runs, so the wiring
+   *    can order a group before calling any of its tasks, and a group holding
+   *    no tasks still exists.

Review Comment:
   Kept the documented ordering and dropped the recursion. It is now 
`createGroup`, a direct parent lookup, since `record` is given parents before 
the groups nested in them.
   
   Fixed in a22621f5749.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt:
##########
@@ -72,6 +72,7 @@ class Bundle(
    */
   fun register(dag: DagDef): Bundle {
     checkOpen()
+    dag.expandGroupEdges()

Review Comment:
   Covered by the change above: `Bundle.register` no longer calls 
`expandGroupEdges()` at all.
   
   Fixed in a22621f5749.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt:
##########
@@ -116,7 +139,40 @@ object Refs {
     }
     inputs.filterIsInstance<TaskRef<*>>().forEach { def.dependsOn(it.def) }
     def.inputs += inputs
-    active.dag.addTask(def)
+    val group = active.groupOf[def.id]

Review Comment:
   Good call, done. The view emits `Refs.node(groupId(), def)` inside a group 
and `Refs.node("", def)` at Dag level. That removed `Recording.groupOf`, the 
`Map.ofEntries(...)` in the builder, the second `record` overload and 
`TaskDeclaration.groupId`.
   
   Fixed in a22621f5749.



##########
java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt:
##########
@@ -326,22 +389,75 @@ class BuilderProcessor : AbstractProcessor() {
   private fun inType(paramType: TypeMirror): TypeName =
     ParameterizedTypeName.get(ARG_TYPE, 
WildcardTypeName.subtypeOf(TypeName.get(paramType).boxIfPossible()))
 
-  private fun collectTasks(el: TypeElement): List<TaskDeclaration> {
-    val declarations = mutableListOf<TaskDeclaration>()
+  /** The Dag's tasks and task groups, read from the class tree the author 
wrote. */
+  private fun collectScope(
+    el: TypeElement,
+    path: List<String>,
+    classPath: List<String>,
+  ): Scope {
+    val tasks = mutableListOf<TaskDeclaration>()
     for (inner in el.enclosedElements) {
       if (inner !is ExecutableElement) continue
       val ann = inner.getAnnotation(Builder.Task::class.java) ?: continue
       if (inner.isVarArgs) throw IllegalArgumentException("Cannot create task 
from vararg function ${inner.simpleName}")
-      val id = ann.id.ifBlank { inner.simpleName.toString() }
-      require(declarations.none { it.id == id }) { "Tasks in Dag have 
duplicate ID: $id" }
-      require(declarations.none { 
it.method.simpleName.contentEquals(inner.simpleName) }) {
-        "Dag class ${el.simpleName} overloads task method 
'${inner.simpleName}'; a method's name is the " +
+      val localId = ann.id.ifBlank { inner.simpleName.toString() }
+      require('.' !in localId) {
+        "Task ID '$localId' on method '${inner.simpleName}' of 
${el.simpleName} contains '.', which " +
+          "Airflow reads as a task group prefix; declare the task inside a 
@Builder.TaskGroup class instead"
+      }
+      require(tasks.none { 
it.method.simpleName.contentEquals(inner.simpleName) }) {
+        "Class ${el.simpleName} overloads task method '${inner.simpleName}'; a 
method's name is the " +
           "name of its generated task class and of its wiring-view method, so 
rename one and keep its " +
-          "task id with @Builder.Task(id = \"$id\")"
+          "task id with @Builder.Task(id = \"$localId\")"
+      }
+      tasks +=
+        TaskDeclaration(
+          inner,
+          (path + localId).joinToString("."),
+          collectDataParams(inner),
+          el,
+          classPath,
+          path.joinToString("."),
+        )
+    }
+
+    val groups = mutableListOf<GroupDeclaration>()
+    for (inner in el.enclosedElements.filterIsInstance<TypeElement>()) {
+      val ann = inner.getAnnotation(Builder.TaskGroup::class.java) ?: continue
+      val localId = checkGroupClass(inner, ann)
+      require(groups.none { it.id == localId }) {
+        "Class ${el.simpleName} declares more than one task group '$localId'"
       }
-      declarations += TaskDeclaration(inner, id, collectDataParams(inner))
+      val scope = collectScope(inner, path + localId, classPath + 
inner.simpleName.toString())
+      groups += GroupDeclaration(inner, localId, (path + 
localId).joinToString("."), scope)
+    }

Review Comment:
   Both added. `checkIds` now rejects a task and a task group sharing an id, 
alongside the duplicate-task check.
   
   For the misplaced class: `process` gained a pass mirroring the 
`@Builder.Deps` one, so a `@Builder.TaskGroup` class not nested in a 
`@Builder.Dag` class or in another `@Builder.TaskGroup` class is an error.
   
   Both fixed in a22621f5749.



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

Reply via email to