← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 65번째

Airflow 모듈 65/151 airflow-learn-65

Go SDK - 실행 아키텍처와 태스크 작성

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/language-sdks/go.rst - Go SDK (도입부), Prerequisites, Execution architecture, Quick start, Writing tasks, sdk.Client surface, Reading the task runtime context, XCom type mapping (약 18-320줄)

이 모듈을 다 읽으면

  • Go SDK에서 Python 태스크 러너가 Go 번들을 어떻게 직접 실행하는지 설명할 수 있다
  • Go 태스크 함수의 파라미터 타입 기반 의존성 주입 방식을 설명할 수 있다
  • Go에서 XCom을 읽을 때 정수 타입을 다뤄야 하는 이유를 설명할 수 있다

Go SDK는 컴파일 언어이므로 모든 태스크를 사전에 컴파일해 Dag 소스와 메타데이터 매니페스트까지 하나의 실행 파일(번들)로 묶어 배포하며, 태스크 함수는 시그니처의 파라미터 타입에 따라 sdk.TIRunContext/*slog.Logger/sdk.Client 등이 자동 주입되는 평범한 Go 함수로 작성한다.

실행 아키텍처 - Python 태스크 러너가 Go 바이너리를 직접 실행

Go SDK를 쓰려면 번들을 빌드/패킹할 때 Go 1.24 이상이 필요하지만, 이는 빌드 시점에만 필요한 요구사항이다 - 패킹된 번들이 실행되는 워커에는 Go 툴체인이 전혀 필요 없다. 번들 자체가 자기 완결적인 네이티브 실행 파일이기 때문이다.

파이썬 태스크 러너가 Go 번들을 직접 실행하며, 호스트에 별도의 Go 워커 프로세스를 띄우지 않는다 - 이는 Java SDK와 동일한 코디네이터 메커니즘이다. 성숙한 파이썬 슈퍼바이저가 Airflow와 맞닿는 관심사를 처리해주기 때문에, Go 태스크는 원격 태스크 로그(S3/GCS), 전체 범위의 태스크 상태, 대체 XCom 백엔드 같은 기능을 Go에서 다시 구현할 필요 없이 그대로 물려받는다.

핵심 포인트

  • Go 툴체인은 빌드 시점에만 필요하고, 패킹된 번들이 실행되는 워커에는 Go가 설치되어 있을 필요가 없다 (자기 완결적 네이티브 실행 파일이기 때문)
  • 파이썬 태스크 러너가 Go 번들을 직접 실행하는 방식이라, 원격 태스크 로그(S3/GCS)나 다양한 태스크 상태, 대체 XCom 백엔드 같은 기능을 Go에서 다시 구현할 필요 없이 성숙한 파이썬 슈퍼바이저의 기능을 그대로 물려받는다

태스크 함수 작성 - 시그니처 기반 의존성 주입

Go 태스크는 평범한 Go 함수이며, 런타임이 함수 시그니처를 살펴 타입에 따라 인자를 주입한다 - 따라서 각 태스크는 자신에게 필요한 파라미터만 선언하면 된다. 주요 주입 대상은 `sdk.TIRunContext`(태스크의 취소/데드라인 신호와 태스크 인스턴스 식별자, Dag 런 타임스탬프), `*slog.Logger`(출력이 Airflow 태스크 로그로 라우팅되는 로거), `sdk.Client`(혹은 그보다 좁은 인터페이스, Variables/Connections/XCom용 클라이언트)다.

선택적인 `(any, error)` 반환값에서 값은 태스크의 `return_value` XCom이 되고, nil이 아닌 `error`(또는 런타임이 복구하는 panic)는 태스크 인스턴스를 실패로 표시해 stub에 설정된 재시도를 트리거한다.

필요한 가장 좁은 인터페이스를 요청하는 것(예: 전체 `sdk.Client` 대신 `sdk.VariableClient`)은 그 태스크가 어떤 Airflow 기능을 쓰는지 문서화하는 효과가 있고, 테스트에서 가짜(fake) 구현을 넣기도 쉽게 만든다. `sdk.Client`는 세 개의 더 작은 인터페이스로 구성된다 - `VariableClient`(`GetVariable`로 문자열 반환, `UnmarshalJSONVariable`로 JSON Variable을 포인터에 디코드), `ConnectionClient`(`GetConnection`이 `ID`, `Type`, `Host`, `Port`, `Login`, `Password`, `Path`, `Extra`(map) 필드와 `GetURI()` 헬퍼를 가진 `Connection`을 반환), `XComClient`(`GetXCom`으로 업스트림 XCom을 읽고 `PushXCom`으로 값을 발행). 값이 없는 조회는 `VariableNotFound`, `ConnectionNotFound`, `XComNotFound` 같은 sentinel 에러를 반환하므로, 에러 문자열을 파싱하는 대신 `errors.Is`로 분기 처리할 수 있다.

핵심 포인트

  • Go 태스크 함수는 시그니처에 선언한 파라미터 타입(sdk.TIRunContext, *slog.Logger, sdk.Client 등)에 따라 런타임이 자동으로 값을 주입하는 평범한 함수이며, 필요한 파라미터만 선언하면 된다
  • 반환값은 (any, error) 형태로, 값은 return_value XCom이 되고 error(혹은 복구된 panic)는 태스크를 실패로 표시해 stub에 설정된 재시도를 트리거한다
  • sdk.Client는 VariableClient/ConnectionClient/XComClient 세 인터페이스로 구성되며, 필요한 것만 좁혀 선언하면 어떤 Airflow 기능을 쓰는지 문서화되고 테스트에서 가짜(fake) 구현을 넣기도 쉬워진다
  • 값이 없는 조회는 VariableNotFound/ConnectionNotFound/XComNotFound 같은 sentinel 에러를 반환하므로 errors.Is로 분기 처리한다

실행 컨텍스트 읽기와 XCom 타입 매핑

`sdk.TIRunContext`는 `context.Context`를 임베드하는 인터페이스이므로, 같은 `ctx`가 취소와 클라이언트 호출을 동시에 담당한다. `ctx.TaskInstance()`는 `DagID`, `RunID`, `TaskID`, `MapIndex`(매핑되지 않은 태스크라면 nil), `TryNumber`를 반환하고, `ctx.DagRun()`은 `DagID`, `RunID`와 `*time.Time` 필드인 `LogicalDate`, `DataIntervalStart`, `DataIntervalEnd`를 반환한다(해당하는 값이 없는 런, 예를 들어 수동 트리거라면 nil).

XCom 값은 Airflow 메타데이터 DB에 JSON으로 저장된다. Python `int`는 JSON 정수로, Go에서는 값에 따라 폭이 달라지는 수치 타입으로 온다. Python `float`는 JSON 소수로, Go `float64`로 온다. `str`→문자열→`string`, `bool`→불리언→`bool`, `None`→null→`nil`, `list`→배열→`[]any`, `dict`→객체→`map[string]any`로 매핑된다.

`GetXCom`은 전송 계층에서 디코딩된 값을 그대로 반환할 뿐, 아직 타입이 지정된 역직렬화 계층이 없다. 파이썬 슈퍼바이저가 값을 msgpack으로 인코딩하기 때문에, 정수값은 그 값에 따라 폭이 달라지는 Go 정수 타입으로 오고 정수가 아닌 값만 `float64`로 온다 - 따라서 고정된 정수 폭을 가정해서는 안 되며, 예상하는 수치 타입들에 대해 타입 스위치를 하거나, `json.Marshal`/`json.Unmarshal`로 원하는 Go 타입에 왕복시켜 안전하게 처리해야 한다.

핵심 포인트

  • ctx.TaskInstance()와 ctx.DagRun()으로 태스크/Dag 런 식별자와 스케줄링 타임스탬프(LogicalDate 등, 수동 트리거 시 nil일 수 있음)를 읽을 수 있다
  • GetXCom이 반환하는 정수는 고정된 폭의 타입이 아니다 - Python 슈퍼바이저가 msgpack으로 인코딩하므로 정수는 값에 따라 폭이 달라지는 Go 정수 타입으로, 정수가 아닌 값만 float64로 온다 - 타입 스위치를 하거나 json.Marshal/Unmarshal로 왕복시켜 안전하게 처리해야 한다