Airflow의 감사 로그

Airflow의 감사 로그 (Audit Logs in Airflow)

감사 로그는 Airflow 시스템의 역사적 기록으로, 누가 어떤 작업을 언제 수행했는지 문서화해요. 이 문서는 감사 로그의 목적, 이벤트 로그와의 차이, 접근 방법, 이벤트 카탈로그와 조회 방법을 설명해 드려요.

출처: 문서

본문

감사 로그 이해하기 (Understanding Audit Logs)

감사 로그는 Airflow 시스템의 역사적 기록 역할을 하며, 누가 어떤 작업을 수행했고 언제 발생했는지 문서화해요. 이 로그는 시스템 무결성 유지, 규정 준수(compliance) 요구 충족, 문제 발생 시 포렌식 분석 수행에 필수적이에요.

본질적으로 감사 로그는 세 가지 기본 질문에 답해요:

  • 누가 (Who): 어떤 사용자나 시스템 컴포넌트가 작업을 시작했는가
  • 무엇 (What): 수행된 특정 작업
  • 언제 (When): 이벤트의 정확한 타임스탬프

감사 로그의 주요 목적은 다음과 같아요:

  • 규제 준수 (Regulatory Compliance): 데이터 거버넌스와 감사 추적 요구 충족
  • 보안 모니터링 (Security Monitoring): 무단 접근이나 의심스러운 활동 감지
  • 운영 문제 해결 (Operational Troubleshooting): 시스템 문제로 이어지는 이벤트의 순서 이해
  • 변경 관리 (Change Management): 중요 시스템 컴포넌트의 수정 추적

Note

감사 로그에 접근하려면 Audit Logs.can_read 권한이 필요해요. 이 권한이 있는 사용자는 DAG별 접근 권한과 관계없이 모든 감사 항목을 볼 수 있어요.

이벤트 로그 이해하기 (Understanding Event Logs)

이벤트 로그는 Airflow 시스템의 운영적 heartbeat를 나타내요. 책임 추적과 규정 준수에 초점을 맞춘 감사 로그와 달리, 이벤트 로그는 시스템 동작, 애플리케이션 성능, 운영 메트릭의 기술적 세부 사항을 포착해요.

이벤트 로그는 몇 가지 중요한 기능을 제공해요:

  • 디버깅 및 문제 해결 (Debugging and Troubleshooting): 상세 오류 메시지와 stack traces 제공
  • 성능 모니터링 (Performance Monitoring): 실행 시간, 리소스 사용량, 시스템 메트릭 기록
  • 운영 통찰 (Operational Insights): 시스템 상태, 컴포넌트 상호작용, 워크플로우 실행 추적
  • 개발 지원 (Development Support): 코드 디버깅과 최적화를 위한 상세 정보 제공

이벤트 로그는 보통 로그 파일이나 외부 로깅 시스템에 저장되며 다음과 같은 정보를 포함해요:

  • Task 실행 세부 정보와 출력
  • 시스템 오류와 경고
  • 성능 메트릭과 타이밍 정보
  • 컴포넌트 시작·종료 이벤트
  • 리소스 사용량 데이터

감사 로그 vs 이벤트 로그

두 로깅 시스템 모두 시스템 관리에 중요하지만, 서로 다른 목적과 대상을 제공해요:

Characteristic Audit Logs Event Logs
Primary Purpose 책임 추적과 규정 준수 추적 운영 모니터링과 시스템 디버깅
Target Audience 보안 팀, 감사자, 규정 준수 담당자 개발자, 시스템 관리자, 운영 팀
Content Focus 사용자 작업과 관리 변경 시스템 동작, 오류, 성능 데이터
Storage Location 구조화된 데이터베이스 테이블 (log) 로그 파일, 외부 로깅 시스템
Retention Requirements 장기 (규정 준수를 위해 수개월~수년, 데이터베이스에서 제거하지 않는다면) 단기중기 (며칠몇 주)
Query Patterns "누가 재실행을 위해 task 인스턴스를 clear했나?" 로그 집계 프레임워크를 사용하지 않으면 질의 없음. 보통 로그는 task 실행별로 읽히며 "이 task 실행이 왜 실패했나?"를 설명

감사 로그 접근하기 (Accessing Audit Logs)

Airflow는 감사 로그 데이터에 접근하는 여러 인터페이스를 제공하며, 각각 다른 사용 사례와 기술 요구에 적합해요:

웹 사용자 인터페이스 (Web User Interface) : Airflow 웹 인터페이스는 감사 로그를 보는 가장 접근하기 쉬운 방법을 제공해요. Browse → Audit Logs로 이동하면 내장 필터링, 정렬, 검색 기능이 있는 인터페이스에 접근할 수 있어요. 이 인터페이스는 임시 조사와 일상 모니터링에 이상적이에요.

REST API 통합 (REST API Integration) : 프로그램적 접근과 시스템 통합을 위해 /eventLogs REST API 엔드포인트를 사용해요. 이 접근 방식은 자동화된 모니터링, 외부 보안 도구와의 통합, 커스텀 보고 애플리케이션을 가능하게 해요.

감사 로깅의 범위 (Scope of Audit Logging)

Airflow의 감사 로깅 시스템은 세 가지 뚜렷한 운영 영역에 걸쳐 이벤트를 포착해요:

사용자 시작 작업 (User-Initiated Actions) : 사용자가 어떤 인터페이스(웹 UI, REST API, 명령줄 도구)를 통해서든 Airflow와 상호작용할 때 발생하는 이벤트예요. 예시:

- 수동 DAG run 트리거와 수정
- variables, connections, pools의 구성 변경
- Task 인스턴스 상태 수정 (clear, success/failed로 mark)
- 관리 작업과 사용자 관리 활동

시스템 생성 이벤트 (System-Generated Events) : 정상 운영 중 Airflow의 내부 프로세스가 자동으로 만드는 이벤트예요:

- Task 라이프사이클 상태 전이 (queued, running, success, failed)
- 시스템 모니터링 이벤트 (heartbeat timeout, 외부 상태 변경)
- 자동 복구 작업 (task 재스케줄, 재시도 시도)
- 리소스 관리 활동

명령줄 인터페이스 작업 (Command-Line Interface Operations) : Airflow CLI 도구를 통해 수행된 활동을 포착하는 이벤트예요:

- 직접 task 실행 명령
- DAG 관리 작업
- 시스템 관리 및 유지보수 작업
- 자동화된 스크립트 실행

일반적인 감사 로그 시나리오

감사 로그 분석을 돕기 위해 자주 접하는 시나리오와 해당 쿼리가 여기 있어요:

"누가 이 DAG을 트리거했나요?"

SELECT dttm, owner, extra
FROM log
WHERE event = 'trigger_dag_run' AND dag_id = 'example_dag'
ORDER BY dttm DESC;

"이 실패한 task에 무슨 일이 있었나요?"

SELECT dttm, event, owner, extra
FROM log
WHERE dag_id = 'example_dag' AND task_id = 'example_task'
ORDER BY dttm DESC;

"최근에 variables를 바꾼 사람은 누구죠?"

SELECT dttm, event, owner, extra
FROM log
WHERE event LIKE '%variable%'
ORDER BY dttm DESC LIMIT 20;

이벤트 카탈로그 (Event Catalog)

다음 섹션은 Airflow의 감사 로깅 시스템이 추적하는 모든 이벤트의 완전한 참조를 제공해요. 이런 이벤트 유형을 이해하면 감사 로그를 해석하고 특정 사용 사례에 대한 효과적인 쿼리를 구성하는 데 도움이 돼요.

Task Instance 이벤트

시스템 생성 task 이벤트:

  • running: Task 인스턴스 실행 시작
  • success: Task 인스턴스 성공 완료
  • failed: Task 인스턴스 실행 중 실패
  • skipped: Task 인스턴스 건너뜀
  • upstream_failed: 업스트림 실패로 Task 인스턴스 실패
  • up_for_retry: Task 인스턴스 재시도 예정
  • up_for_reschedule: Task 인스턴스 재스케줄됨
  • queued: Task 인스턴스 실행을 위해 큐에 대기
  • scheduled: Task 인스턴스 스케줄됨
  • deferred: Task 인스턴스 지연됨 (트리거 대기 중)
  • awaiting_input: Task 인스턴스가 사람 입력 대기 중 (Human-in-the-loop)
  • restarting: Task 인스턴스 재시작 중
  • removed: Task 인스턴스 제거됨

시스템 모니터링 이벤트:

  • heartbeat timeout: Task 인스턴스가 heartbeat 전송을 중단해 종료될 것임
  • state mismatch: Task 인스턴스 상태가 (Airflow 외부에서) 변경됨
  • stuck in queued reschedule: Task 인스턴스가 queued 상태에 갇혀 재스케줄됨
  • stuck in queued tries exceeded: Task 인스턴스가 최대 재큐 시도를 초과함

사용자 시작 task 이벤트:

  • fail task: 사용자가 task를 수동으로 실패로 표시
  • skip task: 사용자가 task를 수동으로 건너뜀으로 표시
  • action_set_failed: 사용자가 UI/API를 통해 task 인스턴스를 실패로 설정
  • action_set_success: 사용자가 UI/API를 통해 task 인스턴스를 성공으로 설정
  • action_set_retry: 사용자가 task 인스턴스를 재시도로 설정
  • action_set_skipped: 사용자가 task 인스턴스를 건너뜀으로 설정
  • action_set_running: 사용자가 task 인스턴스를 실행 중으로 설정
  • action_clear: 사용자가 task 인스턴스 상태를 clear

사용자 액션 이벤트

DAG 작업:

  • trigger_dag_run: 사용자가 DAG run 트리거
  • delete_dag_run: 사용자가 DAG run 삭제
  • patch_dag_run: 사용자가 DAG run 수정
  • clear_dag_run: 사용자가 DAG run clear
  • get_dag_run: 사용자가 DAG run 정보 조회
  • get_dag_runs_batch: 사용자가 여러 DAG runs 조회
  • post_dag_run: 사용자가 DAG run 생성
  • patch_dag: 사용자가 DAG 구성 수정
  • get_dag: 사용자가 DAG 정보 조회
  • get_dags: 사용자가 여러 DAGs 조회
  • delete_dag: 사용자가 DAG 삭제

Task instance 작업:

  • post_clear_task_instances: 사용자가 task instances clear
  • patch_task_instance: 사용자가 task instance 수정
  • get_task_instances_batch: 사용자가 task instance 정보 조회
  • delete_task_instance: 사용자가 task instance 삭제
  • get_task_instance: 사용자가 단일 task instance 정보 조회
  • get_task_instance_tries: 사용자가 task instance 재시도 정보 조회
  • patch_task_instances_batch: 사용자가 여러 task instances 수정

Variable 작업:

  • delete_variable: 사용자가 variable 삭제
  • patch_variable: 사용자가 variable 수정
  • post_variable: 사용자가 variable 생성
  • bulk_variables: 사용자가 일괄 variable 작업 수행

Connection 작업:

  • delete_connection: 사용자가 connection 삭제
  • post_connection: 사용자가 connection 생성
  • patch_connection: 사용자가 connection 수정
  • bulk_connections: 사용자가 일괄 connection 작업 수행
  • create_default_connections: 사용자가 기본 connections 생성

Pool 작업:

  • get_pool: 사용자가 pool 정보 조회
  • get_pools: 사용자가 여러 pools 조회
  • post_pool: 사용자가 pool 생성
  • patch_pool: 사용자가 pool 수정
  • delete_pool: 사용자가 pool 삭제
  • bulk_pools: 사용자가 일괄 pool 작업 수행

Asset 작업:

  • get_asset: 사용자가 asset 정보 조회
  • get_assets: 사용자가 여러 assets 조회
  • get_asset_alias: 사용자가 asset alias 정보 조회
  • get_asset_aliases: 사용자가 여러 asset aliases 조회
  • post_asset_events: 사용자가 asset events 생성
  • get_asset_events: 사용자가 asset events 조회
  • materialize_asset: 사용자가 asset materialization 트리거
  • get_asset_queued_events: 사용자가 queued asset events 조회
  • delete_asset_queued_events: 사용자가 queued asset events 삭제
  • get_dag_asset_queued_events: 사용자가 DAG asset queued events 조회
  • delete_dag_asset_queued_events: 사용자가 DAG asset queued events 삭제
  • get_dag_asset_queued_event: 사용자가 특정 DAG asset queued event 조회
  • delete_dag_asset_queued_event: 사용자가 특정 DAG asset queued event 삭제

Backfill 작업:

  • get_backfill: 사용자가 backfill 정보 조회
  • get_backfills: 사용자가 여러 backfills 조회
  • post_backfill: 사용자가 backfill 생성
  • pause_backfill: 사용자가 backfill 일시정지
  • unpause_backfill: 사용자가 backfill 일시정지 해제
  • cancel_backfill: 사용자가 backfill 취소
  • create_backfill_dry_run: 사용자가 backfill dry run 수행

사용자 및 역할 관리:

  • get_user: 사용자가 사용자 정보 조회
  • get_users: 사용자가 여러 사용자 조회
  • post_user: 사용자가 사용자 계정 생성
  • patch_user: 사용자가 사용자 계정 수정
  • delete_user: 사용자가 사용자 계정 삭제
  • get_role: 사용자가 역할 정보 조회
  • get_roles: 사용자가 여러 역할 조회
  • post_role: 사용자가 역할 생성
  • patch_role: 사용자가 역할 수정
  • delete_role: 사용자가 역할 삭제

CLI 이벤트

DAG 관리 명령:

  • cli_dags_list: 시스템의 모든 DAGs 나열
  • cli_dags_show: DAG 정보와 구조 표시
  • cli_dags_state: DAG run 상태 확인
  • cli_dags_next_execution: DAG의 다음 실행 시간 표시
  • cli_dags_trigger: 명령줄에서 DAG run 트리거
  • cli_dags_delete: DAG과 그 메타데이터 삭제
  • cli_dags_pause: DAG 일시정지
  • cli_dags_unpause: DAG 일시정지 해제
  • cli_dags_backfill: 날짜 범위에 대해 DAG runs 백필
  • cli_dags_test: 데이터베이스에 영향 없이 DAG 테스트

Task 관리 명령:

  • cli_tasks_list: 특정 DAG의 tasks 나열
  • cli_tasks_run: 특정 task instance 실행
  • cli_tasks_test: 데이터베이스에 영향 없이 task 테스트
  • cli_tasks_state: task instance 상태 확인
  • cli_tasks_failed_deps: task의 실패한 의존성 표시
  • cli_tasks_render: task 템플릿 렌더링
  • cli_tasks_clear: task instance 상태 clear

데이터베이스 및 시스템 명령:

  • cli_db_init: Airflow 데이터베이스 초기화
  • cli_db_upgrade: 데이터베이스 스키마 업그레이드
  • cli_db_reset: 데이터베이스 재설정 (위험한 작업)
  • cli_db_shell: 데이터베이스 셸 열기
  • cli_db_check: 데이터베이스 연결성과 스키마 확인
  • cli_db_migrate: 데이터베이스 스키마 마이그레이션 (레거시 명령)
  • cli_migratedb: 레거시 데이터베이스 마이그레이션 명령
  • cli_initdb: 레거시 데이터베이스 초기화 명령
  • cli_resetdb: 레거시 데이터베이스 재설정 명령
  • cli_upgradedb: 레거시 데이터베이스 업그레이드 명령

사용자 및 보안 명령:

  • cli_users_create: 새 사용자 계정 생성
  • cli_users_delete: 사용자 계정 삭제
  • cli_users_list: 모든 사용자 나열
  • cli_users_add_role: 사용자에게 역할 추가
  • cli_users_remove_role: 사용자에게서 역할 제거

구성 및 Variable 명령:

  • cli_variables_get: variable 값 조회
  • cli_variables_set: variable 값 설정
  • cli_variables_delete: variable 삭제
  • cli_variables_list: 모든 variables 나열
  • cli_variables_import: 파일에서 variables 가져오기
  • cli_variables_export: variables를 파일로 내보내기

Connection 관리 명령:

  • cli_connections_get: connection 세부 정보 조회
  • cli_connections_add: 새 connection 추가
  • cli_connections_delete: connection 삭제
  • cli_connections_list: 모든 connections 나열
  • cli_connections_import: 파일에서 connections 가져오기
  • cli_connections_export: connections를 파일로 내보내기

Pool 관리 명령:

  • cli_pools_get: pool 정보 얻기
  • cli_pools_set: pool 생성 또는 업데이트
  • cli_pools_delete: pool 삭제
  • cli_pools_list: 모든 pools 나열
  • cli_pools_import: 파일에서 pools 가져오기
  • cli_pools_export: pools를 파일로 내보내기

서비스 및 프로세스 명령:

  • cli_webserver: Airflow webserver 시작
  • cli_scheduler: Airflow scheduler 시작
  • cli_worker: Celery worker 시작
  • cli_flower: Flower 모니터링 도구 시작
  • cli_triggerer: triggerer 프로세스 시작
  • cli_standalone: Airflow를 standalone 모드로 시작
  • cli_api_server: Airflow API server 시작
  • cli_dag_processor: DAG processor 서비스 시작
  • cli_celery_worker: Celery worker 시작 (대체 명령)
  • cli_celery_flower: Celery Flower 시작 (대체 명령)

유지보수 및 유틸리티 명령:

  • cli_cheat_sheet: CLI 명령 참조 표시
  • cli_version: Airflow 버전 정보 표시
  • cli_info: 시스템 정보 표시
  • cli_config_get_value: 구성 값 가져오기
  • cli_config_list: 구성 옵션 나열
  • cli_plugins: 설치된 플러그인 나열
  • cli_rotate_fernet_key: Fernet 암호화 키 회전
  • cli_sync_perm: 권한 동기화
  • cli_shell: 대화형 Python 셸 시작
  • cli_kerberos: Kerberos 티켓 renewer 시작

테스트 및 개발 명령:

  • cli_test: 테스트 실행
  • cli_render: 템플릿 렌더링
  • cli_dag_deps: DAG 의존성 표시
  • cli_task_deps: task 의존성 표시

레거시 명령:

  • cli_run: 레거시 task run 명령
  • cli_backfill: 레거시 backfill 명령
  • cli_clear: 레거시 clear 명령
  • cli_list_dags: 레거시 DAG list 명령
  • cli_list_tasks: 레거시 task list 명령
  • cli_pause: 레거시 pause 명령
  • cli_unpause: 레거시 unpause 명령
  • cli_trigger_dag: 레거시 DAG trigger 명령

각 CLI 명령 감사 로그 항목은 다음을 포함해요:

  • 사용자 식별 (User identification): 명령을 실행한 사람
  • 명령 세부 정보 (Command details): 인자를 포함한 전체 명령
  • 실행 컨텍스트 (Execution context): 작업 디렉터리, 환경 변수
  • 타임스탬프 (Timestamp): 명령이 실행된 시각
  • 종료 상태 (Exit status): 성공 또는 실패 표시

커스텀 이벤트 (Custom Events)

Airflow는 커스텀 감사 로그 항목을 프로그래밍 방식으로 만들 수 있게 해줘요:

from airflow.models.log import Log
from airflow.utils.session import provide_session


@provide_session
def log_custom_event(session=None):
    log_entry = Log(event="custom_event", owner="username", extra="Additional context information")
    session.add(log_entry)
    session.commit()

감사 로그 항목의 구조 (Anatomy of an Audit Log Entry)

각 감사 로그 레코드는 로그된 이벤트의 완전한 그림을 제공하는 구조화된 정보를 포함해요. 이 필드를 이해하는 것은 효과적인 로그 분석에 필수적이에요:

Field Name Description and Usage
dttm 이벤트가 발생한 시각을 나타내는 타임스탬프 (UTC 시간대)
event 작업 또는 이벤트의 설명적 이름 (예: trigger_dag_run, failed)
owner 행위자의 신원: 사용자 작업의 경우 username, 시스템 이벤트의 경우 "airflow"
dag_id 영향받은 DAG의 식별자 (해당하는 경우)
task_id 영향받은 task의 식별자 (해당하는 경우)
run_id 실행 인스턴스 추적을 위한 특정 DAG run 식별자
try_number task 재시도 및 재실행을 위한 시도 번호
map_index 동적으로 매핑된 tasks의 인덱스
logical_date DAG run의 논리적 실행 날짜
extra JSON 형식의 추가 컨텍스트 (파라미터, 오류 세부 정보 등)

감사 로그 조회 방법 (Audit Log Query Methods)

효과적인 감사 로그 분석은 로그 데이터를 조회하고 검색하기 위해 사용 가능한 다양한 방법을 이해하는 것을 요구해요. 각 방법은 장점이 있으며 다른 시나리오에 적합해요:

REST API 예시:

# Get all audit logs
curl -X GET "http://localhost:8080/api/v1/eventLogs"

# Filter by event type
curl -X GET "http://localhost:8080/api/v1/eventLogs?event=trigger_dag_run"

# Filter by DAG
curl -X GET "http://localhost:8080/api/v1/eventLogs?dag_id=example_dag"

# Filter by date range
curl -X GET "http://localhost:8080/api/v1/eventLogs?after=2024-01-01T00:00:00Z&before=2024-12-31T23:59:59Z"

데이터베이스 쿼리 예시:

-- Get recent user actions
SELECT dttm, event, owner, dag_id, task_id, extra
FROM log
WHERE owner IS NOT NULL
ORDER BY dttm DESC
LIMIT 100;

-- Get task failure events
SELECT dttm, dag_id, task_id, run_id, extra
FROM log
WHERE event = 'failed'
ORDER BY dttm DESC;

-- Get user actions on specific DAG
SELECT dttm, event, owner, extra
FROM log
WHERE dag_id = 'example_dag' AND owner IS NOT NULL
ORDER BY dttm DESC;

이벤트 로그 조회하기 (Querying Event Logs)

이벤트 로그(운영 로그)는 보통 로깅 구성에 따라 다른 방법으로 접근돼요:

로그 파일:

# View scheduler logs
tail -f $AIRFLOW_HOME/logs/scheduler/latest/*.log

# View webserver logs
tail -f $AIRFLOW_HOME/logs/webserver/webserver.log

# View task logs for specific DAG run
cat $AIRFLOW_HOME/logs/dag_id/task_id/2024-01-01T00:00:00+00:00/1.log

Task 로그용 REST API:

# Get task instance logs
curl -X GET "http://localhost:8080/api/v1/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}/logs/{try_number}"

# Get task logs with metadata
curl -X GET "http://localhost:8080/api/v1/dags/example_dag/dagRuns/2024-01-01T00:00:00+00:00/taskInstances/example_task/logs/1?full_content=true"

Python 로깅 통합:

import logging
from airflow.utils.log.logging_mixin import LoggingMixin


class MyOperator(BaseOperator, LoggingMixin):
    def execute(self, context):
        # These will appear in event logs
        self.log.info("Task started")
        self.log.warning("Warning message")
        self.log.error("Error occurred")

외부 로깅 시스템:

외부 로깅 시스템(예: ELK stack, Splunk, CloudWatch)를 사용할 때:

# Example Elasticsearch query
curl -X GET "elasticsearch:9200/airflow-*/_search" -H 'Content-Type: application/json' -d'
{
  "query": {
    "bool": {
      "must": [
        {"match": {"dag_id": "example_dag"}},
        {"range": {"@timestamp": {"gte": "2024-01-01", "lte": "2024-01-31"}}}
      ]
    }
  }
}'

실용적인 쿼리 예시 (Practical Query Examples)

다음 예시는 일반적인 운영 및 보안 시나리오에 대한 감사 로그 쿼리의 실용적 응용을 보여줘요. 이 쿼리는 특정 요구사항에 맞게 조정할 수 있는 템플릿 역할을 해요:

보안 조사 (Security Investigation)

-- Find all actions by a specific user in the last 24 hours
SELECT dttm, event, dag_id, task_id, extra
FROM log
WHERE owner = 'suspicious_user'
  AND dttm > NOW() - INTERVAL '24 hours'
ORDER BY dttm DESC;

규정 준수 보고 (Compliance Reporting)

-- Get all variable and connection changes for audit report
SELECT dttm, event, owner, extra
FROM log
WHERE event IN ('post_variable', 'patch_variable', 'delete_variable',
                'post_connection', 'patch_connection', 'delete_connection')
  AND dttm BETWEEN '2024-01-01' AND '2024-01-31'
ORDER BY dttm;

DAG 문제 해결 (Troubleshooting DAG Issues)

-- See all events for a problematic DAG run
SELECT dttm, event, task_id, owner, extra
FROM log
WHERE dag_id = 'example_dag'
  AND run_id = '2024-01-15T10:00:00+00:00'
ORDER BY dttm;

더 알아보기 (Learn more)