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 clearget_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 clearpatch_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;