uranusjr commented on code in PR #74230:
URL: https://github.com/apache/airflow/pull/74230#discussion_r4203855231
##########
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:
We should also check there’s not a task and a task group using the same name
(which is not allowed in core).
--
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]