jason810496 commented on code in PR #73895: URL: https://github.com/apache/airflow/pull/73895#discussion_r4140123509
########## go-sdk/airflow/inputs.go: ########## @@ -0,0 +1,143 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package airflow + +import ( + "fmt" + "reflect" + "slices" + "strconv" + "strings" +) + +type inputs []*TaskRef + +func (in inputs) applyTask(c *taskConfig) { c.inputs = append(c.inputs, in) } + +// Inputs passes the results of tasks to the task that [DagRef.Task] adds, and makes each of +// those tasks an upstream task of the new one. It is the Go form of a Python TaskFlow call such +// as transform(extract()): +// +// extracted := dag.Task(extract) +// transformed := dag.Task(transform, airflow.Inputs(extracted)) +// dag.Task(load, airflow.Inputs(transformed)) +// +// The result of the first task fills the first parameter after the [Context], the result of the +// second task fills the second parameter, and so on. For the Dag above, the task functions could +// be: +// +// func extract(actx airflow.Context) ([]string, error) +// func transform(actx airflow.Context, rows []string) (int, error) +// func load(actx airflow.Context, count int) error +// +// A struct parameter takes the whole result, even when it is the only parameter after the +// Context. A function passed to [TaskHandler] binds the fields of such a struct by name instead. Review Comment: The statement here is kind of incorrect: It should be something like if the `airflow.Inputs()` only receive one struct, we will bind them by name no matter it's TaskHandler or a Task of a Dag. ########## go-sdk/airflow/task_option.go: ########## @@ -17,17 +17,20 @@ package airflow -// TaskOption is an option to [DagRef.Task]. [TaskSpec] is one. +// TaskOption is an option to [DagRef.Task]. There are two kinds: a [TaskSpec] sets the +// attributes of the task that DagRef.Task adds, and [Inputs] passes the results of other tasks +// to that task. // // Its only method is unexported, so a type outside this package cannot declare it. // A struct that embeds a TaskSpec or a TaskOption still satisfies the interface, and Task // panics when it is given one. type TaskOption interface{ applyTask(*taskConfig) } type taskConfig struct { - // specs keeps every TaskSpec passed to DagRef.Task, so that Task can reject a second one - // instead of merging the two. - specs []TaskSpec + // specs and inputs keep every TaskSpec and every Inputs passed to DagRef.Task, so that Task + // can reject a second TaskSpec or a second Inputs instead of merging it into the first. + specs []TaskSpec + inputs [][]*TaskRef Review Comment: Not necessary in this PR, but a follow-up after we complete the feature gap of native Dag that would it make more sense to have a bool to represent the `inputs` is stored or not and directly reject the second incoming inputs at the `applyTask` stage instead of storing all the `inputs` down. So does the `specs`. Not urgent and no need to address this as current approach works. Tracked in https://github.com/orgs/apache/projects/499/views/7?pane=issue&itemId=258243802. ########## go-sdk/airflow/dag.go: ########## @@ -163,19 +177,32 @@ func (d *DagRef) Task(fn any, opts ...TaskOption) *TaskRef { )) } } - if _, exists := d.taskIDs[taskID]; exists { + if _, exists := d.tasksByID[taskID]; exists { panic(fmt.Sprintf( "airflow.DagRef.Task: Dag %q already has a task %q; "+ "set another task_id with airflow.TaskSpec{TaskID: ...}", d.dagID, taskID, )) } - if d.taskIDs == nil { - d.taskIDs = make(map[string]struct{}) + fnType := reflect.TypeOf(fn) + upstreams := d.checkInputs(taskID, fnType, cfg.inputs) + var resultType reflect.Type + if fnType.NumOut() == 2 { Review Comment: Would it be better to raise error if there're more or equal than three the return arguments? -- 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]
