Airflow - Astro CLI

조정훈·2024년 4월 8일

Airflow 구성요소

  1. Webserver: Flask로 구현된 Airflow UI 서버

  2. Scheduler: 작업을 스케줄링하는 파이썬 멀티스레딩 프로세스, 어떤 작업을 언제 어디에서 수행할지 결정

  3. Database: DAG, Task들의 메타데이터를 저장하는 데이터베이스, 일반적으로 PostgreSQL을 사용

  4. Executor: Task를 실행하는 메커니즘, Airflow가 작동할 때 마다 스케줄러에서 실행

  5. Worker: 실제 Task를 수행하는 프로세스, 어떤 Executor를 선택하느냐에 따라 해당 Worker를 인프라에 가지고 있을수도 아닐 수도 있다. (로컬 Worker를 사용하거나, 외부 Worker를 사용할 수 있다는 의미)



Astro CLI 설치

Astro CLI는 Airflow를 개발한 오픈소스 팀 Astronomer가 개발환경을 구축하기 쉽게 만들어둔 프로그램
도커를 미리 설치해 놓자 (docker 설치 관련 글 있으니 확인)

설치명령어 : curl -sSL install.astronomer.io | sudo bash -s -- v1.25.0

초기설정 : astro dev init 하면 프로젝트 파일에 폴더들이 생긴다.

git branch로 main으로 이동해주고 pull 한번 땡긴다.

그 후 -> 레파지토리 가서 issue 추가

ex)  [FEAT] Airflow 초기 세팅

	background 
		airflow에 필요한 초기 세팅을 진행합니다.
	todo
		[ ] astro dev init 파일 추가

이슈 생성했으니 브랜치 만들기 git checkout -b feat-3/airflow-init

생성된 파일들 다 추가 git add . (평상시에는 . 쓰지말고 이렇게 최초 세팅시에만 .을 사용하는 습관을 들이자)

커밋

git commit -m 'feat - airflow init
>
> #3'  (이슈번호 맞게 적어주자)

dags 경로에 파이프라인을 작성하면 된다. dags 경로에 추가된 .py파일을 인식하여 파이프라인이 추가된다.


최초로 한번 실행하기 : astro dev start
이 명령어를 진행하면 로컬의 postgre db를 사용하게 되어있는데
나는 클라우드의 postgresql 을 사용할거라서 이 부분을 삭제해 줄 것이다.

실행 됐으면 astro dev kill 명령어로 생성해놓은것을 지운다.

그 후 astro dev object export --compose 해서 dags 관련 파일을 도커컴포즈 파일로 내보내준다.

그 후 좌측 파일에 compose.yaml 파일에서 수정해줘야한다. (postgres 관련 삭제)

삭제할부분:

아래로 내려보면서 depends_on 포함해서 postgres 들어가부분 다 지우기 (postgres만 있으면 depends_on 까지 같이 삭제, 다른것도 있으면 postgres만 삭제)

ctrl + F 해서 postgres 치면 위에 주소 부분만 나와야함.

끝났으면

astro dev start --compose-file compose.yaml

이 명령어로 지정된 Docker Compose 파일에 정의된 서비스를 사용하여 Airflow를 시작한다. 이 명령어를 사용하면 특정 환경에 맞게 커스터마이징된 Docker Compose 파일을 사용하여 Airflow를 시작할 수 있다.

기다리면 Airflow is starting up! 메세지가 나오는데 좀 더 기다려보면

실행되는데 제대로 됐으면 db에 table 제대로 들어간걸 확인해 볼 수 있다.(deaver)


그리고 주소창에 loaclhost:8080 치면 airflow 접속할 수 있다.

<로그인> id: admin / pw: admin


자주 쓰는 기능 설명

DAGs

DAG을 clear 하고 싶다 -> bar를 클릭해서 clear task 하면 되고

task를 clear하고 싶다 -> 작은 칸을 눌러서 clear task 하면 된다.

clear task 하게 되면 해당 작업을 삭제시키는데, airflow에서 자체적으로 마치 시간을 되돌려서 과거로 돌아가 task를 다시 시작하는것 과 같은 상태를 만들어준다. 오류가 생겼거나, 결과를 무효화 하고 싶거나, DAG 구조를 변경 후 테스트하고 싶거나, 다시 실행하고 싶을 때 사용된다.

  • 주의해야 할 점

    Clear Task를 사용하면 이미 완료된 Task의 결과도 모두 삭제됩니다.
    Clear Task는 특정 Task만 선택적으로 실행할 수 없으며, 해당 Task와 그 하위 Task들이 모두 Clear됩니다.
    Clear Task를 사용하기 전에 DAG 구조와 Task 실행 순서를 충분히 이해하고 있어야 합니다.


Admin

Connections

-> 민감한 정보들을 변수로 지정해놓고 쓸려고 설정 하는 곳


Pools

-> 동시성 제어 설정 하는 곳. 동시에 시행될 수 있는 리소스들을 정의해놓는다.
리소스 제한/ 우선 순위 부여 / 병렬 실행 제한 등을 할 수 있다.

ex1) database pool 해놓으면 커넥션의 제한이 있고 database나 쿼리엔진에 순간적인 부하를 주면 서로 영향을 주게된다. 그래서 Bigquery 혹은 PostgreSQL 같은 SQL연산 하는 작업이 있다면 pool 의 갯수를 제한을 해놓는다. 그러면 pool에 속한 애들이 동시에 돌 수 있는 최대 수(concurrency:동시성)를 제한한다.

ex2) 여러 테스트 파이프라인을 실행하는 Airflow 환경이 있다고 가정. 그러나 각 테스트 파이프라인은 외부 API 호출을 통해 데이터를 가져와야 한다. 외부 API는 동시에 여러 요청을 처리할 수 있지만, 무제한으로 요청을 보내면 서버에 부하를 일으킬 수 있따. 이때 Pools를 사용하여 외부 API 호출에 대한 최대 동시 요청 수를 제한할 수 있다

XCom



참고링크) cron 표현식 정리

0개의 댓글