FrankYang0529 commented on code in PR #73767: URL: https://github.com/apache/airflow/pull/73767#discussion_r4117653449
########## go-sdk/airflow/dag.go: ########## @@ -0,0 +1,222 @@ +// 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" + "runtime" + "strings" + "sync" +) + +// DagRef is a Dag authored in Go. [Dag] returns a new one. +type DagRef struct { + dagId string + // Dag and Task copy the specs they are given, so a caller cannot change a registered Dag + // through a spec it still holds. A slice or map field in DagSpec or TaskSpec would share its + // contents with the caller, so Dag and Task would have to copy that field too. + spec DagSpec + + mu sync.Mutex + registered bool + tasks []*TaskRef + taskIds map[string]struct{} +} + +// Dag returns an empty Dag with the given dag_id. spec holds the rest of the Dag's attributes. +// Add the tasks with [DagRef.Task], then pass the Dag to [BundleRef.Register]: +// +// dag := airflow.Dag("etl", airflow.DagSpec{}) +// dag.Task(extract) +// dag.Task(load, airflow.TaskSpec{TaskID: "load_rows"}) +// +// bundle.Register(dag) +// +// Add every task before Register. [DagRef.Task] panics once the Dag is registered. +// +// [BundleRef.Serve] does not yet serve the Dags that Dag returns. It leaves them out of the +// --airflow-metadata manifest and cannot run their tasks. +func Dag(dagId string, spec DagSpec) *DagRef { Review Comment: Changed `Dag` to `func Dag(dagId string, spec ...DagSpec) *DagRef`, so `airflow.Dag("etl")` works. Passing more than one `DagSpec` panics, as passing a second `TaskSpec` to `dag.Task` does. I also updated the signature in ADR 0008. -- 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]
