커맨드라인 인터페이스(CLI)로 Airflow 다루기
커맨드라인 인터페이스(CLI)로 Airflow 다루기
Airflow의 CLI는 DAG 구조를 이미지로 뽑는 것부터 데이터베이스 정리, 커넥션 내보내기까지 운영 중에 자주 쓰는 작업들을 손쉽게 처리해줘요. 웹 UI가 부담스러운 순간이나 자동화 스크립트 안에서든, CLI 하나로 대부분의 일을 해낼 수 있어요.
출처: Using the Command Line Interface — Airflow Documentation
본문
Bash/Zsh 자동 완성 설정
bash(또는 zsh)를 쓴다면 airflow는 argcomplete로 자동 완성을 제공해요.
argcomplete를 쓰는 모든 Python 애플리케이션에서 전역 활성화를 하려면:
sudo activate-global-python-argcomplete
전역이 아닌 Airflow 영구 활성화에는:
register-python-argcomplete airflow >> ~/.bashrc
Airflow만을 위한 1회성 활성화에는:
eval "$(register-python-argcomplete airflow)"
zsh를 쓴다면 .zshrc에 다음을 추가해요:
autoload bashcompinit
bashcompinit
eval "$(register-python-argcomplete airflow)"
커넥션 생성
CLI로 커넥션을 만드는 방법은 Creating a Connection from the CLI에서 자세히 볼 수 있어요.
DAG 구조를 이미지로 내보내기
Airflow는 DAG 구조를 이미지로 출력하거나 저장할 수 있어요. 문서화하거나 공유할 때 유용해요. Graphviz가 설치되어 있어야 해요.
예를 들어 example_complex DAG를 터미널에 출력하려면:
airflow dags show example_complex
이러면 렌더링된 DAG 구조가 DOT 형식으로 화면에 출력돼요.
여러 파일 포맷이 지원돼요. --save[filename].[format] 인자를 붙이면 됩니다.
example_complex DAG를 PNG로 저장하려면:
airflow dags show example_complex --save example_complex.png
지원되는 파일 포맷은 bmp, dot, gv, xdot, pdf, png, svg, jpg, jpeg, eps, json 등 아주 다양해요.
기본적으로 Airflow는 airflow.cfg의 [core] 섹션에 있는 dags_folder 옵션으로 지정된 디렉터리에서 DAG를 찾아요.
--subdir 인자는 Airflow 3.x에서 더 이상 지원되지 않아요. 특정 파일의 DAG를 테스트하려면:
airflow dags test <DAG_ID> -f <path_to_dag_file>
DAG 구조 표시하기
복잡한 의존성을 가진 DAG를 다룰 때는 DAG를 미리 봐서 올바른지 확인하는 게 좋아요.
macOS에서는 iTerm2와 imgcat 스크립트를 함께 쓰면 콘솔에 DAG 구조를 표시할 수 있어요. Graphviz도 설치되어 있어야 해요.
airflow dags show 명령에 --imgcat 스위치를 쓰면 돼요. 예를 들어 example_bash_operator DAG를 표시하려면:
airflow dags show example_bash_operator --imgcat
명령 출력 서식 바꾸기
airflow dags list나 airflow tasks states-for-dag-run 같은 일부 명령은 --output 플래그로 출력 서식을 바꿀 수 있어요. 가능한 옵션은 다음과 같아요.
table: 정보를 일반 텍스트 테이블로simple: 표준 Linux 유틸리티로 파싱 가능한 단순 테이블로json: json 문자열 형태로yaml: 유효한 yaml 형태로
json과 yaml 형식은 jq나 yq 같은 커맨드라인 도구로 데이터를 다루기 쉬워요. 예를 들어:
airflow tasks states-for-dag-run example_complex 2020-11-13T00:00:00+00:00 --output json | jq ".[] | {sd: .start_date, ed: .end_date}"
{ "sd": "2020-11-29T14:53:46.811030+00:00", "ed": "2020-11-29T14:53:46.974545+00:00" }
{ "sd": "2020-11-29T14:53:56.926441+00:00", "ed": "2020-11-29T14:53:57.118781+00:00" }
{ "sd": "2020-11-29T14:53:56.915802+00:00", "ed": "2020-11-29T14:53:57.125230+00:00" }
메타데이터 데이터베이스에서 기록 정리하기
db clean 명령을 실행하기 전에 메타데이터 데이터베이스를 백업하는 것을 강력히 권장해요.
db clean 명령은 각 테이블에서 지정한 --clean-before-timestamp보다 오래된 레코드를 삭제하는 방식으로 동작해요. 삭제할 테이블 목록을 선택적으로 지정할 수 있고, 목록을 주지 않으면 모든 테이블이 포함돼요.
--dag-ids(쉼표 구분 목록)로 특정 DAG만, 또는 --exclude-dag-ids(쉼표 구분 목록)로 특정 DAG를 제외하고 정리를 필터링할 수 있어요. 이 옵션들로 특정 DAG에 대한 정리를 타겟하거나 회피할 수 있어요.
--dry-run 옵션을 쓰면 정리 대상 기본 테이블의 행 수를 출력해줘요.
기본적으로 db clean은 정리된 행을 _airflow_deleted__<table>__<timestamp> 형태의 테이블에 아카이브해요. 이렇게 보존하고 싶지 않다면 --skip-archive 인자를 주세요.
--skip-archive 없이 실행하다 에러가 나면 _airflow_deleted__<table>__<timestamp> 테이블이 DB에 남아요. db drop-archived 명령으로 이 테이블들을 수동으로 제거할 수 있어요.
정리 실패 감지
기본적으로 db clean은 테이블별 오류(예: 매우 큰 테이블에서 statement_timeout 초과)를 억제하고, 일부 테이블이 정리되지 않았어도 코드 0으로 종료해요. 건너뛴 테이블 목록은 로그에 WARNING으로 남아요.
어떤 테이블 정리가 실패하면 0이 아닌 코드로 종료시켜, DAG 태스크에서 airflow db clean을 호출할 때 실패 시 태스크가 빨간불이 되게 하려면 --error-on-cleanup-failure을 주세요.
airflow db clean \
--clean-before-timestamp "$(date -u -d '21 days ago' '+%Y-%m-%dT%H:%M:%S+00:00')" \
--yes \
--error-on-cleanup-failure
--error-on-cleanup-failure가 설정되면, 발생한 RuntimeError가 정리에 실패한 테이블 목록을 포함해 어떤 테이블이 정리되지 않았는지 알려줘요.
아카이브 CREATE TABLE … AS SELECT 단계 자체가 타임아웃되기 쉬운 대규모 배포에서는 --error-on-cleanup-failure와 --skip-archive를 함께 쓰는 것을 권장해요. --skip-archive는 중간 아카이브 테이블 없이 행을 직접 삭제해서, 연산이 더 빠르고 statement_timeout에 걸릴 확률이 낮아요.
아카이브 테이블에서 정리된 기록 내보내기
db export-archived 명령은 db clean이 만든 아카이브 테이블의 내용을 지정한 형식으로, 기본적으로 CSV로 내보내요. 내보낸 파일에는 db clean 과정에서 기본 테이블에서 정리된 레코드가 담겨요.
--export-format 옵션으로 형식을 지정할 수 있는데 기본 형식은 csv이며 현재 유일하게 지원되는 형식이에요. 내보낼 경로는 --output-path 옵션으로 지정해야 하고, 해당 위치는 존재해야 해요.
그 외 옵션으로는 내보낼 테이블을 지정하는 --tables, 내보낸 뒤 아카이브 테이블을 삭제하는 --drop-archives가 있어요.
아카이브 테이블 삭제하기
db clean 과정에서 아카이브 테이블을 삭제하는 --skip-archive 옵션을 쓰지 않았다면, 여전히 db drop-archived 명령으로 아카이브 테이블을 삭제할 수 있어요. 이 연산은 되돌릴 수 없으므로, 삭제 전에 db export-archived 명령으로 테이블을 디스크에 백업하는 걸 권장해요.
--tables 옵션으로 삭제할 테이블을 지정할 수 있고, 지정하지 않으면 모든 아카이브 테이블이 삭제돼요.
연쇄 삭제 주의
일부 테이블은 ON DELETE CASCADE로 외래 키 관계가 정의되어 있어서, 한 테이블의 삭제가 다른 테이블의 삭제를 촉발할 수 있어요. 예를 들어 task_instance 테이블은 dag_run 테이블을 키로 참조하므로, DagRun 레코드가 삭제되면 연관된 모든 태스크 인스턴스도 함께 삭제돼요.
DAG run 특별 처리
Airflow는 보통 최신 DagRun을 조회해 다음에 실행할 DagRun을 결정해요. 모든 DAG run을 삭제하면, catchup=True가 설정된 경우 이미 완료된 옛 DAG run을 스케줄할 수도 있어요. 그래서 db clean은 스케줄링 연속성을 위해 가장 최근의 수동 트리거가 아닌 DagRun을 보존해요.
backfill 가능한 DAG 고려사항
모든 DAG가 backfill 명령과 함께 쓰도록 설계된 것은 아니에요. 백필 가능한 DAG라면 주의가 필요해요. DAG runs를 삭제하고, 삭제된 DAG runs를 포함하는 날짜 범위에 대해 백필을 실행하면 그 runs가 다시 만들어져 실행돼요. 이런 범주에 속하는 DAG는 DAG runs 삭제를 자제하고, 태스크 인스턴스나 로그 같은 큰 테이블만 정리하는 게 좋아요.
Airflow 업그레이드
airflow db migrate --help로 사용법을 확인할 수 있어요.
업그레이드 마이그레이션 수동 실행
원한다면 업그레이드용 SQL 문을 생성해서 각 업그레이드 마이그레이션을 하나씩 수동으로 적용할 수 있어요. db migrate에 --range(Airflow 버전용)나 --revision-range(Alembic 리비전용) 옵션을 쓰면 돼요. Alembic 리비전 id 업데이트 명령을 건너뛰지 마세요. 이것이 다음 업그레이드 때 어디서부터인지 Airflow가 알게 하는 방법이에요. 버전과 리비전의 매핑은 Database Migrations Reference를 보세요.
Airflow 다운그레이드
db downgrade나 다른 데이터베이스 연산 전에는 데이터베이스를 백업하는 것을 권장해요.
db downgrade 명령으로 특정 Airflow 버전으로 다운그레이드할 수 있고, Alembic 리비전 id로 대신 지정할 수도 있어요. 명령을 미리 보기만 하고 실행하지 않으려면 --show-sql-only 옵션을 쓰세요.
--from-revision과 --from-version 옵션은 --show-sql-only 옵션과 함께만 쓸 수 있어요. 실제 마이그레이션을 실행할 때는 항상 현재 리비전에서 다운그레이드해야 하기 때문이에요.
Airflow 다운그레이드(Python 환경에 설치된 Airflow 버전을 내린 직후, 데이터베이스를 내린 직후가 아님)를 마친 뒤에는 dags reserialize로 DAG를 다시 직렬화하는 것을 강력히 권장해요. 직렬화된 DAG가 다운그레이드된 버전과 호환되도록 하기 위함이에요.
커넥션 내보내기
CLI로 데이터베이스에서 커넥션을 내보낼 수 있어요. 지원되는 파일 형식은 json, yaml, env예요.
대상 파일을 파라미터로 지정할 수 있어요:
airflow connections export connections.json
file-format 파라미터로 파일 형식을 덮어쓸 수도 있어요:
airflow connections export /tmp/connections --file-format yaml
-를 지정하면 STDOUT으로 출력해요:
airflow connections export -
JSON 형식은 키가 커넥션 id, 값이 커넥션 정의인 객체를 담아요. 커넥션은 JSON 객체로 정의돼요. 다음은 JSON 파일 예시예요.
{
"airflow_db": {
"conn_type": "mysql",
"host": "mysql",
"login": "root",
"password": "plainpassword",
"schema": "airflow",
"port": null,
"extra": null
},
"druid_broker_default": {
"conn_type": "druid",
"host": "druid-broker",
"login": null,
"password": null,
"schema": null,
"port": 8082,
"extra": "{\"endpoint\": \"druid/v2/sql\"}"
}
}
YAML 파일 구조도 JSON과 비슷해요. 커넥션 id와 하나 이상의 커넥션 정의 키-값 쌍으로 구성되며, 커넥션은 YAML 객체로 정의돼요. 다음은 YAML 파일 예시예요.
airflow_db:
conn_type: mysql
extra: null
host: mysql
login: root
password: plainpassword
port: null
schema: airflow
druid_broker_default:
conn_type: druid
extra: '{"endpoint": "druid/v2/sql"}'
host: druid-broker
login: null
password: null
port: 8082
schema: null
.env 형식으로도 내보낼 수 있어요. 키는 커넥션 id이고, 값은 Airflow의 Connection URI 형식이나 JSON으로 직렬화된 커넥션 표현이에요. JSON을 쓰려면 --serialization-format=json 옵션을 주고, 아니면 Airflow Connection URI 형식이 사용돼요. 두 형식의 .env 파일 예시는 다음과 같아요.
URI 예:
airflow_db=mysql://root:plainpassword@mysql/airflow
druid_broker_default=druid://druid-broker:8082?endpoint=druid%2Fv2%2Fsql
JSON 출력 예:
airflow_db={"conn_type": "mysql", "login": "root", "password": "plainpassword", "host": "mysql", "schema": "airflow"}
druid_broker_default={"conn_type": "druid", "host": "druid-broker", "port": 8082, "extra": "{\"endpoint\": \"druid/v2/sql\"}"}
DAG import 에러 테스트
CLI는 list-import-errors 하위 명령으로 발견된 DAG에 import 에러가 있는지 확인할 수 있어요. 커맨드 출력을 확인해 import되지 못하는 DAG가 있으면 실패시키는 자동화 단계를 만들 수 있어요. 특히 --output으로 표준 형식을 지정할 때 유용해요. 에러가 없을 때의 기본 출력은 No data found이고, json 출력은 []예요. 이 체크를 CI나 pre-commit 훅에서 돌려 리뷰 과정과 테스트를 빠르게 할 수 있어요.
jq로 출력을 파싱해 에러가 있으면 실패하는 명령 예시:
airflow dags list-import-errors --output=json | jq -e 'select(type=="array" and length == 0)'
이 줄은 자동화에 그대로 넣을 수 있고, 출력도 보고 싶다면 tee를 쓸 수 있어요:
airflow dags list-import-errors | tee import_errors.txt && jq -e 'select(type=="array" and length == 0)' import_errors.txt
Jenkins 파이프라인의 예:
stage('All Dags are loadable') {
steps {
sh 'airflow dags list-import-errors | tee import_errors.txt && jq -e \'select(type=="array" and length == 0)\' import_errors.txt'
}
}
정확하게 동작하려면 Airflow가 stdout에 추가 텍스트를 로그로 남기지 않도록 해야 해요. 예를 들어 deprecation 경고를 고치거나, 로드될 때 로그를 만드는 플러그인이 있다면 lazy_load_plugins=True를 설정하거나, 2>/dev/null을 추가해야 할 수 있어요.