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]

Reply via email to