DAG 파일 처리

DAG 파일 처리 (Dag File Processing)

이 페이지는 DAG를 정의하는 Python 파일을 읽어서 Scheduler가 스케줄링할 수 있도록 저장하는 DAG 파일 처리(Dag File Processing) 과정을 설명해요. DagFileProcessorManager가 어떤 파일을 처리할지 결정하고, DagFileProcessorProcess가 개별 파일을 하나 이상의 DAG 객체로 변환해요. 마지막에는 Dag processor 성능을 미세 조정하는 방법과 설정 옵션도 다뤄요.

출처: 문서

본문

Dag File Processing은 DAG를 정의하는 Python 파일을 읽고, Scheduler가 스케줄링할 수 있도록 저장하는 과정을 가리켜요.

Dag 파일 처리에는 두 가지 주요 구성 요소가 있어요. DagFileProcessorManager는 어떤 파일을 처리해야 할지 결정하는 무한 루프를 실행하는 프로세스이고, DagFileProcessorProcess는 개별 파일을 하나 이상의 DAG 객체로 변환하기 위해 시작되는 별도의 프로세스예요.

DagFileProcessorManager는 사용자 코드를 실행해요. 그 결과 이 프로세스는 airflow dag-processor CLI 명령을 실행해 독립 프로세스로 실행돼요.

../_images/dag_file_processing_diagram.png

DagFileProcessorManager는 다음과 같은 단계를 거쳐요:

  1. 새 파일 확인: DAG가 마지막으로 새로고침된 이후 경과 시간이 refresh_interval보다 크면 파일 경로 목록을 업데이트해요.
  2. 최근 처리 파일 제외: min_file_process_interval보다 더 최근에 처리됐고 수정되지 않은 파일을 제외해요.
  3. 파일 경로 대기열에 넣기: 발견한 파일을 파일 경로 대기열에 추가해요.
  4. 파일 처리: 각 파일에 대해 최대 parsing_processes개까지 새 DagFileProcessorProcess를 시작해요.
  5. 결과 수집: 완료된 Dag processor로부터 결과를 수집해요.
  6. 통계 기록: 통계를 출력하고 dag_processing.total_parse_time을 내보내요.

DagFileProcessorProcess는 다음과 같은 단계를 거쳐요:

  1. 파일 처리: 전체 프로세스는 dag_file_processor_timeout 안에 완료되어야 해요.
  2. DAG 파일을 Python 모듈로 로드: dagbag_import_timeout 안에 완료되어야 해요.
  3. 모듈 처리: Python 모듈 안에서 DAG 객체를 찾아요.
  4. DagBag 반환: DagFileProcessorManager에 발견된 DAG 객체 목록을 제공해요.

Dag processor 성능 미세 조정

Dag processor 성능에 영향을 주는 것들

Dag processor는 DAG 파일을 계속 파싱하고 DB의 DAG와 동기화하는 책임을 져요. Dag processor를 미세 조정하려면 여러 요소를 고려해야 해요:

  • 나의 배포 종류
    • DAG를 공유하기 위해 어떤 종류의 파일시스템을 쓰는지 (DAG를 계속 읽는 성능에 영향을 줌)
    • 파일시스템이 얼마나 빠른지 (분산 클라우드 파일시스템의 대부분은 더 많은 비용을 지불해 더 높은 처리량/빠른 파일시스템을 얻을 수 있음)
    • 처리에 사용할 메모리가 얼마나 있는지
    • 사용 가능한 CPU가 얼마나 있는지
    • 사용 가능한 네트워킹 대역폭이 얼마나 있는지
  • Dag 구조의 로직과 정의:
    • DAG 파일이 몇 개인지
    • 파일 안에 DAG가 몇 개인지
    • DAG 파일이 얼마나 큰지 (Dag 파서가 매 n초마다 파일을 읽고 파싱해야 한다는 점을 기억하세요)
    • DAG가 얼마나 복잡한지 (얼마나 빨리 파싱될 수 있는지, Task·의존성이 몇 개인지)
    • DAG 파일을 파싱하는 것이 많은 라이브러리를 import하거나 최상위(toplevel)에서 무거운 처리를 하는지 (힌트! 그러면 안 됩니다. Top level Python Code 참고)
  • Dag processor 설정
    • Dag processor가 몇 개인지
    • Dag processor에 parsing process가 몇 개인지
    • Dag processor가 같은 DAG를 다시 파싱하기 전에 기다리는 시간 (계속 일어남)
    • Dag processor 루프마다 실행하는 콜백 수

Dag processor 미세 조정 접근법

Airflow는 성능을 미세 조정하기 위해 많은 "knob(손잡이)"을 제공하지만, 어떤 knob을 어떻게 돌려 최상의 효과를 얻을지는 사용자의 특정 배포, DAG 구조, 하드웨어 가용성, 기대치에 따라 달라지는 별개의 작업이에요. 배포를 관리하는 일의 일부는 무엇을 최적화할지 결정하는 것이에요. 어떤 사용자는 CPU 사용을 낮추는 대가로 새 DAG 파싱에 30초 지연을 허용하는 반면, 다른 사용자는 예를 들어 더 높은 CPU 사용을 감수하면서 Dags 폴더에 나타난 DAG가 거의 즉시 파싱되길 기대해요.

Airflow는 결정할 유연성을 주지만, 어떤 성능 측면이 자신에게 가장 중요한지 파악하고 어떤 knob을 어느 방향으로 돌릴지 결정해야 해요.

일반적으로 미세 조정을 위한 접근 방식은 어떤 성능 개선·최적화와도 동일해야 해요 (특정 도구를 추천하진 않을게요 — 시스템을 관찰·모니터링할 때 평소 사용하는 도구를 쓰면 돼요):

  • 평소 시스템을 모니터링할 때 쓰는 올바른 도구 세트로 시스템을 모니터링하는 것이 매우 중요해요. 이 문서는 사용할 수 있는 특정 메트릭·도구에 대해 자세히 다루지 않고, 어떤 종류의 자원을 모니터링해야 하는지만 설명해요. 올바른 데이터를 얻기 위해서는 모니터링에 대한 모범 사례를 따르면 돼요.
  • 어떤 성능 측면이 자신에게 가장 중요한지 결정해요 (무엇을 개선하고 싶은지).
  • 시스템을 관찰해 병목이 어디인지 확인해요: CPU, 메모리, I/O가 일반적인 제한 요소예요.
  • 기대치와 관찰에 기반해 다음 개선이 무엇인지 결정하고, 다시 성능·병목 관찰로 돌아가요. 성능 개선은 반복적인 과정이에요.

Dag processor 성능을 제한할 수 있는 자원

주의해야 할 자원 사용 영역이 몇 가지 있어요:

  • 파일시스템 성능. Airflow Dag processor는 (때로는 많은) Python 파일 파싱에 크게 의존하며, 이런 파일은 종종 공유 파일시스템에 있어요. Dag processor는 그 파일을 계속 읽고 다시 파싱해요. 같은 파일을 worker가 사용할 수 있어야 하므로, 종종 분산 파일시스템에 저장되죠. 이를 위해 다양한 파일시스템을 사용할 수 있어요(NFS, CIFS, EFS, GCS fuse, Azure File System이 좋은 예). 이런 파일시스템에는 제어할 수 있는 여러 파라미터가 있고 성능을 미세 조정할 수 있지만, 이는 이 문서의 범위를 벗어나요. 문제가 파일시스템 성능에서 오는지 파악하기 위해 파일시스템의 통계·사용량을 관찰해야 해요. 예를 들어, EFS를 사용할 때 EFS 성능을 위한 IOPS를 늘리는(그리고 더 많은 비용을 지불하는) 것이 Airflow DAG 파싱의 안정성과 속도를 극적으로 향상시킨다는 일화적 증거가 있어요.
  • 파일시스템 성능이 병목이 된다면, DAG를 배포하는 대안적인 매커니즘으로 전환하는 것이 또 다른 해결책이에요. 이미지에 DAG를 내장하는 방법과 GitSync 배포는 모두 파일이 Dag processor에 로컬로 제공된다는 속성이 있어요. 그래서 분산 파일시스템을 사용해 파일을 읽을 필요가 없어요. 파일이 로컬로 제공되므로, 특히 로컬 저장소에 빠른 SSD 디스크를 사용한다면 보통 가능한 한 빠르게 처리돼요. 이런 배포 매커니즘에는 사용자에게 최선의 선택이 아닐 수 있는 다른 특성도 있지만, 성능 문제가 분산 파일시스템 성능 때문에 발생한다면 이들이 최고의 접근법일 수 있어요.
  • 성능을 높이고 더 많은 것을 병렬로 처리하려 할 때 데이터베이스 연결과 데이터베이스 사용이 문제가 될 수 있어요. Airflow는 "데이터베이스 연결을 많이 소비하는(DB-connection hungry)" 것으로 알려져 있어요 — DAG가 많을수록, 병렬로 많이 처리할수록 더 많은 데이터베이스 연결이 열려요. 이는 보통 MySQL에서는 문제가 되지 않는데, MySQL은 연결 처리가 스레드 기반이기 때문이에요. 하지만 Postgres에서는 문제가 될 수 있는데, 연결 처리가 프로세스 기반이기 때문이에요. 중간 규모의 Postgres 기반 Airflow 설치도 있다면 데이터베이스 프록시로 PGBouncer를 사용하는 것이 최선이라는 게 일반적인 합의예요. Apache Airflow용 Helm Chart는 PGBouncer를 기본적으로 지원해요.
  • CPU 사용은 FileProcessor에게 가장 중요해요 — Python Dag 파일을 파싱·실행하는 프로세스들이죠. Dag processor가 보통 이런 파싱을 계속 트리거하므로, DAG가 많을 때 처리가 많은 CPU를 차지할 수 있어요. min_file_process_interval을 늘려 완화할 수 있지만, 이것은 앞서 말한 트레이드오프 중 하나라 파일 변경이 더 느리게 반영되고, 파일을 제출한 뒤 Airflow UI에서 볼 수 있고 Scheduler가 실행하기까지 지연이 생겨요. DAG를 구축하는 방식을 최적화하고 외부 데이터 소스를 피하는 것이 CPU 사용을 개선하는 최선의 접근법이에요. 더 많은 CPU가 있다면 처리 스레드 수 parsing_processes를 늘릴 수 있어요.
  • Airflow는 더 많은 성능을 얻으려 할 때 상당히 많은 메모리를 사용할 수 있어요. Airflow에서 더 많은 성능은 종종 로드를 처리하는 프로세스 수를 늘려 얻어지는데, 각 프로세스는 전체 Python 인터프리터가 로드되고 많은 클래스가 import되며 임시 인메모리 저장소가 필요해요. Airflow는 forking과 copy-on-write 메모리를 사용해 많은 부분을 최적화하지만, forking 후 새 클래스가 import되면 추가 메모리 압력이 생길 수 있어요. 시스템이 가진 것보다 더 많은 메모리를 사용하고 있는지 관찰해야 해요 — 그러면 swap 디스크를 사용하게 되고, 이는 성능을 극적으로 떨어뜨려요. 메모리 사용량을 볼 때는 관찰하는 메모리의 종류에 주의하세요. 보통 total memory used보다는 working memory를 봐야 해요 (이름은 배포에 따라 다를 수 있음).

Dag processor 성능을 개선하려면

자원 사용량을 알게 되면 고려할 수 있는 개선은 다음과 같아요:

  • DAG 최상위 Python 코드의 로직·파싱 효율을 개선하고 복잡성을 줄이기. 이 코드는 계속 파싱되므로 그 코드를 최적화하면 엄청난 개선을 가져올 수 있어요. DAG를 파싱하면서 외부 데이터베이스 등에 접근하려고 하는 경우 특히 그렇죠(이것은 절대 피해야 해요). Top level Python Code는 최상위 Python 코드 작성의 모범 사례를 설명해요. Dag 복잡성 줄이기 문서는 코드 복잡성을 줄이고 싶을 때 살펴볼 수 있는 영역을 제공해요.
  • 자원 활용률을 개선하기. 시스템에 충분히 활용되지 않는 여유 용량이 있을 때(다시 말해 CPU, 메모리 I/O, 네트워킹이 주요 후보) parsing process 수를 늘리는 같은 조치가 더 높은 활용을 대가로 성능 개선을 가져올 수 있어요.
  • 하드웨어 용량을 늘리기 (예: CPU가 제한이거나, DAG 파일시스템에 사용하는 I/O가 한계일 때). 종종 Dag processor 성능 문제는 단순히 시스템이 "충분히 유능"하지 않아서이며, 공유 데이터베이스나 파일시스템이 병목이 아니라면 이것이 유일한 방법일 수 있어요.
  • "Dag processor 튜너블(tunables)" 값으로 다양한 실험해 보기. 단순히 한 성능 측면을 다른 것과 교환해 더 나은 효과를 얻는 경우가 많아요. 예를 들어 CPU 사용을 줄이고 싶다면 파일 처리 간격을 늘릴 수 있어요 (하지만 결과적으로 새 DAG가 더 큰 지연으로 나타나요). 보통 성능 튜닝은 서로 다른 측면을 균형 잡는 예술이에요.
  • 때로는 특정 배포에 더 나은 미세 조정 결과를 얻기 위해 Dag processor 동작을 약간 바꾸기도 해요 (예: 파싱 정렬 순서 변경).

Dag processor 설정 옵션

다음 설정을 사용해 Dag processor의 다양한 측면을 제어할 수 있어요. 그러나 [dag_processor] 섹션에서 사용할 수 있는 다른 비-성능 관련 Dag processor 설정 파라미터도 Configuration Reference에서 볼 수 있어요.

  • file_parsing_sort_mode Dag processor는 파싱 순서를 결정하기 위해 Dag 파일을 나열하고 정렬해요.
  • min_file_process_interval DAG 파일이 다시 파싱된 후 몇 초가 지나면 재파싱하는지를 나타내는 값이에요. DAG 파일은 min_file_process_interval초마다 파싱돼요. DAG에 대한 업데이트는 이 간격 후에 반영돼요. 이 숫자를 낮게 유지하면 CPU 사용이 증가해요.
  • parsing_processes Dag processor는 DAG 파일을 파싱하기 위해 여러 프로세스를 병렬로 실행할 수 있어요. 이 옵션은 몇 개의 프로세스가 실행될지 정의해요.

더 알아보기 (Learn more)