Go SDK
Go SDK
이 페이지는 Airflow Task 로직을 Go로 구현할 수 있게 해주는 Go SDK(실험적 기능)를 다뤄요. DAG와 스케줄링은 Python에 남고, 개별 Task는 각 task instance마다 ExecutableCoordinator가 실행하는 컴파일된 Go bundle에 위임해요. Task 작성, XCom 타입 매핑, 빌드·패키징, 배포, 제한 사항을 설명해요.
출처: 문서
본문
이것은 실험적 기능이에요.
Go SDK는 Airflow의 "모델"(Variables, Connections, XCom)에 대한 네이티브 접근과 함께 Airflow Task 로직을 Go로 구현할 수 있게 해줘요. DAG와 스케줄링은 Python에 남고, 개별 Task는 각 task instance마다 ExecutableCoordinator가 실행하는 컴파일된 Go bundle에 위임해요.
Go는 컴파일 언어이므로, 모든 Task는 미리 컴파일되어 번들이라는 단일·자립형 네이티브 실행 파일 안에 등록되어야 해요. 번들은 또한 실행 파일에 추가된 footer에 자체 Dag 소스와 메타데이터 매니페스트(dag_id·task_id 맵)를 내장하므로, 실행 파일 자체가 번들이에요: 배포할 실행 가능한 단일 파일이고, 별도 매니페스트나 아카이브가 없어요. airflow-go-pack 도구가 그 번들을 빌드·패킹해요.
API 참조
Go SDK 모듈에 대한 생성된 API 참조와 릴리스된 버전 목록은 pkg.go.dev에서 볼 수 있어요.
선행 조건
- 번들을 빌드·패킹하려면 Go 1.24 이상. 이는 빌드 시점 요구사항일 뿐이며, 패킹된 번들을 실행하는 worker는 Go 도구 체인이 필요 없어요. 번들이 자립형 네이티브 실행 파일이니까요.
- 패킹된 번들은 coordinator가 스캔하는 디렉터리 아래에서 Airflow worker가 접근할 수 있어야 해요.
apache-airflow-task-sdk패키지(Airflow와 함께 설치)가 coordinator를 제공해요. 추가 Python 패키지는 필요 없어요.
배포 모드
패킹된 번들은 두 가지 방식으로 실행될 수 있어요. 같은 바이너리가 둘 다에서 작동하며, 배포별로 하나를 선택해요:
- Coordinator (권장). Python Task 러너가 Go 번들을 직접 실행하며, 호스트에 별도 Go worker 프로세스가 없어요. 이는 Java SDK가 사용하는 것과 같은 coordinator 매커니즘이에요. 성숙한 Python supervisor가 Airflow-facing 우려를 처리하므로, 이 경로는 원격 Task 로그(S3/GCS), 전체 Task 상태 범위, 대체 XCom 백엔드를 새로 구현하지 않고 상속해요. 이는 정확히 Edge Worker 경로에 여전히 없는 기능들이에요.
- Edge Worker. 장기 실행 Go 프로세스(
airflow-go-edge-worker)가 Airflow에서 작업을 폴링하고 번들을 실행하며, 데이터 경로에 Python이 없어요. 오늘날 end-to-end로 실행되지만 Limitations에 나열된 기능이 빠져 있어요.
나머지 가이드는 권장하는 coordinator 경로를 다룹니다. Edge Worker 요약은 대안: Go Edge Worker를 참고해요.
빠른 시작
다음 예제는 최소한의 움직이는 부분을 보여줘요: 두 개의 stub task를 가진 Python Dag와 그 Task들의 Go 구현.
Python DAG (스케줄링 쪽)
from airflow.sdk import dag, task
@dag
def simple_dag():
@task.stub(queue="golang")
def extract(): ...
@task.stub(queue="golang")
def transform(): ...
extract() >> transform()
simple_dag()
@task.stub은 Python 구현 없이 Go Task의 모양(이름과 의존성)을 선언해요. queue 값이 Task를 Go coordinator로 라우팅해요.
Go 구현
Task는 일반 Go 함수예요. 런타임이 그 시그니처를 검사하고 타입별로 인자를 주입하므로, 각 Task는 필요한 파라미터만 선언해요.
import (
"log/slog"
"runtime"
"github.com/apache/airflow/go-sdk/sdk"
)
func extract(ctx sdk.TIRunContext, client sdk.Client, log *slog.Logger) (any, error) {
conn, err := client.GetConnection(ctx, "test_http")
if err != nil {
return nil, err
}
log.Info("fetched connection", "host", conn.Host)
// ... do work, honour ctx cancellation ...
return map[string]any{"go_version": runtime.Version()}, nil
}
func transform(ctx sdk.TIRunContext, client sdk.VariableClient, log *slog.Logger) error {
val, err := client.GetVariable(ctx, "my_variable")
if err != nil {
return err
}
log.Info("obtained variable", "my_variable", val)
return nil
}
Note
다른 언어 SDK와 마찬가지로 XCom 의존성은 Python stub Dag에서 선언돼요(그것들이 Task 순서를 정의하니까요). 값은 여전히
client.GetXCom으로 Go에서 명시적으로 읽어야 하며, Task의(any, error)반환값이나client.PushXCom으로 생성돼요.
Go 진입점
bundlev1.BundleProvider를 구현해 Dags와 tasks를 등록해요. main은 한 줄이에요. RegisterDags는 이 번들이 실행할 수 있는 dag_id·Task 이름에 대한 단일 진실의 원천이므로, 생성된 매니페스트가 바이너리가 실제로 실행하는 것과 결코 어긋날 수 없어요.
import (
"log"
v1 "github.com/apache/airflow/go-sdk/bundle/bundlev1"
"github.com/apache/airflow/go-sdk/bundle/bundlev1/bundlev1server"
)
type myBundle struct{}
var _ v1.BundleProvider = (*myBundle)(nil)
func (m *myBundle) GetBundleVersion() v1.BundleInfo {
return v1.BundleInfo{Name: bundleName, Version: &bundleVersion}
}
func (m *myBundle) RegisterDags(dagbag v1.Registry) error {
simpleDag := dagbag.AddDag("simple_dag") // must match the Python dag_id
simpleDag.AddTask(extract) // task_id is taken from the function name
simpleDag.AddTask(transform)
return nil
}
func main() {
if err := bundlev1server.Serve(&myBundle{}); err != nil {
log.Fatal(err)
}
}
AddDag에 전달된 dag_id는 Python DAG의 dag_id와 일치해야 하고, 등록된 각 Task의 이름은 그 DAG의 @task.stub 함수와 일치해야 해요.
Coordinator 구성
airflow.cfg의 [sdk] 아래(또는 동등한 AIRFLOW__SDK__* 환경 변수)에서 coordinator를 등록하고 queue를 라우팅해요:
[sdk]
coordinators = {
"go": {
"classpath": "airflow.sdk.coordinators.executable.ExecutableCoordinator",
"kwargs": {"executables_root": ["~/airflow/executable-bundles"]}
}
}
queue_to_coordinator = {"golang": "go"}
executables_root는 coordinator가 번들을 스캔하는 하나 이상의 디렉터리이고, queue_to_coordinator는 queue="golang"인 stub task를 이 Go coordinator로 라우팅해요. 허용된 kwargs 전체 목록은 ExecutableCoordinator 구성을 참고해요.
실행할 별도 Go worker는 없어요. Airflow worker가 task instance당 번들 바이너리를 한 번 fork해요.
Note
Coordinator는 Airflow worker의 일부이므로,
[sdk]설정(그리고executables_root의 번들 파일)은 Task가 실제로 실행되는 곳에만 있으면 돼요.CeleryExecutor에서는 Celery workers에 설정하는 것으로 충분해요.LocalExecutor에서는 Task가 scheduler 프로세스 안에서 실행되므로, scheduler가 읽을 수 있는 곳에 설정해야 해요. API server와 Dag processor는 필요 없어요.
Task 작성
런타임은 Task 함수의 시그니처를 검사하고 타입별로 인자를 주입하므로, 실제로 필요한 파라미터만 선언하면 돼요:
| 파라미터 타입 | 주입된 값 |
|---|---|
sdk.TIRunContext |
Task의 실행 컨텍스트: 취소/데드라인 신호와 task instance 식별자·Dag run 타임스탬프. 장기 작업에서는 이를 존중하세요. Task 런타임 컨텍스트 읽기 참고. |
*slog.Logger |
출력이 Airflow Task 로그로 라우팅되는 로거. |
sdk.Client (또는 더 좁은 인터페이스) |
Airflow Variables, Connections, XCom용 클라이언트. |
선택적 (any, error) 반환값이 Task의 return_value XCom이 돼요. nil이 아닌 error(또는 런타임이 복구하는 panic)는 Airflow에서 task instance를 실패로 표시하고, stub에 재시도가 구성되어 있으면 재시도를 트리거해요.
필요한 가장 좁은 인터페이스를 요청하면(예: 전체 sdk.Client 대신 sdk.VariableClient) Task가 어떤 Airflow 기능을 건드리는지 문서화되고, 테스트에서 fake를 전달할 수 있어 단위 테스트가 쉬워져요.
sdk.Client 표면
sdk.Client는 세 개의 더 작은 인터페이스를 구성하므로, Task는 그중 하나에만 의존할 수 있어요:
VariableClient—GetVariable(Variable을 문자열로 반환)과UnmarshalJSONVariable(JSON Variable을 제공한 포인터로 디코드).ConnectionClient—Connection을 반환하는GetConnection. 필드ID,Type,Host,Port,Login,Password,Path,Extra(map[string]any)와GetURI()헬퍼가 있어요.XComClient— 상위 Task의 XCom을 읽는GetXCom과 하나를 게시하는PushXCom.
GetXCom은 저장된 값을 any로 반환해요. 저장된 JSON이 Go 타입으로 어떻게 매핑되는지는 XCom 타입 매핑을 참고해요.
찾을 수 없는 조회는 센티널 에러 — VariableNotFound, ConnectionNotFound, XComNotFound — 를 반환하므로, 에러 문자열을 파싱하는 대신 errors.Is로 누락된 값을 분기할 수 있어요.
Task 런타임 컨텍스트 읽기
Task에 sdk.TIRunContext 파라미터를 선언하면 실행 중인 task instance와 그 Dag run의 식별자·스케줄링 타임스탬프를 읽을 수 있어요 — Python·Java SDK가 노출하는 실행 컨텍스트의 Go 버전이에요. context.Context를 내장하는 인터페이스라서, 같은 ctx가 취소와 클라이언트 호출을 구동해요. 런타임은 다른 주입된 파라미터처럼 타입으로 바인딩해요:
func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
ti := ctx.TaskInstance()
log.Info("running",
"dag_id", ti.DagID,
"run_id", ti.RunID,
"task_id", ti.TaskID,
"try_number", ti.TryNumber,
"logical_date", ctx.DagRun().LogicalDate,
)
return nil, nil
}
ctx.TaskInstance()는 DagID, RunID, TaskID, MapIndex(매핑되지 않은 Task는 nil), TryNumber를 반환해요. ctx.DagRun()은 DagID, RunID, 그리고 *time.Time 필드 LogicalDate, DataIntervalStart, DataIntervalEnd(run에 그런 값이 없으면 nil, 예: 수동 트리거)를 반환해요.
XCom 타입 매핑
XCom 값은 Airflow 메타데이터 데이터베이스에 JSON으로 저장돼요. 아래 표는 GetXCom으로 다시 읽을 때 그 JSON 타입이 Go 값으로 어떻게 나타나는지 보여줘요.
| Python 타입 | JSON | Go 타입 (GetXCom으로부터) |
|---|---|---|
int |
number (integer) | numeric (주석 참고) |
float |
number (decimal) | float64 |
str |
string | string |
bool |
boolean | bool |
None |
null | nil |
list |
array | []any |
dict |
object | map[string]any |
Note
GetXCom은 전송에서 디코드된 그대로 값을 반환해요. 아직 타입 있는 XCom 역직렬화 계층은 없어요. 따라서 numeric 값의 구체적 타입은 배포 모드에 따라 달라져요. Execution API를 통해(Edge Worker 경로) 숫자는encoding/json으로 디코드되므로 모든 숫자 — 정수든 아니든 —float64로 도착해요. Coordinator 모드에서는 Python supervisor가 값을msgpack으로 재인코드하므로, 정수는 Go 정수 타입(너비는 값에 따라 달라짐)으로 도착하고 비정수만float64예요. 고정 numeric 타입을 가정하지 마세요: 기대하는 numeric 타입들에 대해 타입 스위치를 하거나,json.Marshal/json.Unmarshal을 통해 값을 타입 있는 Go 값으로 왕복시키세요.
빌드·패키징
일반 go build는 실행 가능한 바이너리를 만들지만, 배포 가능한 번들(바이너리 + 내장 소스 + 매니페스트)은 airflow-go-pack으로 만들어야 해요. 패커는 번들을 컴파일하고 내장 메타데이터 footer를 추가해서, coordinator가 바이너리를 실행하지 않고도 그 dag_id들을 읽고, 단일 실행 파일을 만들 수 있어요. 패커가 내보내는 온디스크 형식(AFBNDL01 footer와 airflow-metadata.yaml 매니페스트)은 모든 네이티브 실행 SDK가 공유하는 번들 형식이며, Executable Bundle Spec에 명시돼 있어요.
airflow-go-pack은 Go 1.24 tool 지시문을 통해 제공되므로 전역 설치가 없어요. 다음을 추가해요:
tool github.com/apache/airflow/go-sdk/cmd/airflow-go-pack
번들 모듈의 go.mod에 추가하고 go tool airflow-go-pack으로 실행해요. 이는 프로젝트별로 패커 버전을 고정해요.
빌드와 패킹을 한 단계로 해요. -- 이후의 플래그는 go build에 그대로 전달돼요:
go tool airflow-go-pack ./example/bundle -- -trimpath -tags=prod
--output <path>를 사용해 패킹된 번들을 coordinator가 스캔하는 디렉터리(executables_root)에 직접 써요:
go tool airflow-go-pack --output ~/airflow/executable-bundles/sample-dag-bundle ./example/bundle
크로스 플랫폼 빌드
번들을 실행하는 worker는 보통 빌드 머신과 다른 운영 체제나 CPU 아키텍처를 사용해요(예: Apple silicon darwin/arm64 랩톱에서 Linux 호스트로 배포). --goos / --goarch를 전달하면 패커가 크로스 빌드해요:
go tool airflow-go-pack --goos linux --goarch amd64 \
--output ~/airflow/executable-bundles/sample-dag-bundle \
./example/bundle
또는 --executable / --source로 미리 빌드된 바이너리를 패킹해요. 패커는 보통 --airflow-metadata로 바이너리를 실행해 매니페스트를 읽지만, 크로스 컴파일된 바이너리는 빌드 호스트에서 실행될 수 없어요. 그 경우 바이너리를 실행할 수 있는 머신에서 매니페스트를 생성하고 --airflow-metadata로 패커에 공급해요:
# On a linux/amd64 machine:
go build -o my-bundle ./example/bundle
./my-bundle --airflow-metadata > airflow-metadata.yaml
# Back on the darwin/arm64 machine:
go tool airflow-go-pack --executable ./my-bundle --source main.go \
--airflow-metadata airflow-metadata.yaml
(--executable은 --goos/--goarch와 -- 이후의 go build 플래그와 상호 배타적이에요. 이미 빌드된 바이너리를 빌드하지 않고 패킹하니까요.)
배포
패킹된 번들을 coordinator의 executables_root에 나열된 디렉터리로 복사하거나 마운트해요. ExecutableCoordinator는 그 디렉터리를 재귀적으로 스캔하고, 들어오는 dag_id를 각 번들의 매니페스트와 일치시키며, 번들의 무결성 해시를 검증하고 일치하는 번들을 실행해요. 번들은 trailer 매직으로 식별되며 파일 이름으로 식별되지 않으므로(리눅스/macOS에선 확장자 없음, Windows에선 .exe), worker의 파일 이름은 무관해요.
ExecutableCoordinator 구성
coordinators 설정 항목의 모든 kwargs는 ExecutableCoordinator 생성자로 전달돼요:
| 파라미터 | 기본값 | 설명 |
|---|---|---|
executables_root |
(필수) | 실행 가능한 번들을 재귀적으로 스캔하는 하나 이상의 디렉터리. 문자열, 경로, 또는 문자열/경로 리스트를 허용해요. |
task_startup_timeout |
10.0 |
실행 후 번들 서브프로세스가 연결될 때까지 기다리는 초. 번들 시작이 느리면(예: 제한된 하드웨어) 이 값을 늘려요. |
대안: Go Edge Worker
같은 번들 바이너리는 Python coordinator 없이도, 독립형 Go Edge Worker 아래에서 실행될 수 있어요. worker가 task당 한 번 번들을 실행하는 대신, airflow-go-edge-worker는 scheduler에 등록하고 Edge Executor API에서 워크로드를 폴링하며, HashiCorp go-plugin(gRPC)을 통해 데이터 경로에 Python 없이 번들을 직접 실행하는 장기 실행 Go 프로세스예요. 번들 소스와 RegisterDags 등록은 동일하고, 배포만 다르며, 모드는 CLI 플래그에서 시작 시 자동 선택되므로 Task 코드를 바꿀 필요가 없어요.
이 경로는 [sdk] coordinators 설정을 사용하지 않으며 현재 Limitations에 나열된 기능이 빠져 있어요. Edge Worker 설정(airflow-go-edge-worker 구성과 go install)은 Go SDK의 자체 저장소 문서를 참고해요.
제한 사항
- Python stub DAG가 여전히 필요해요. Execution API가 아직 비-Python 언어에 대한 DAG 구조를 전달하지 않으므로, Task 이름과 의존성이
@task.stub으로 Python에서 선언돼요. 이는 두 배포 모드 모두에 적용되며 문서화된 알려진 제한 사항이에요.
다음은 Edge Worker 경로가 아직 구현하지 않은 기능의 비-완전 목록이에요. 이는 coordinator 경로가 권장되는 주된 이유예요: coordinator 모드에서는 Python supervisor가 이 우려들을 처리하므로, 거기서는 제한 사항이 아니에요.
- Task를 success나 failed/up-for-retry 이외의 상태로 넣기(deferred, failed-without-retries 등).
- 원격 Task 로그(예: S3/GCS).
- 기본이 아닌 XCom 백엔드를 통한 XCom 읽기/쓰기.