lakeFS와 Apache Airflow 함께 사용하기
Apache Airflow는 사용자가 워크플로를 프로그래밍 방식으로 작성하고, 예약하고, 모니터링할 수 있게 해 주는 플랫폼이에요.
Airflow를 lakeFS와 함께 실행하려면 몇 단계를 따라야 해요.
출처: 문서
본문
Airflow에 lakeFS 연결 만들기
lakeFS 서버에 접근하고 인증하려면 HTTP 타입의 새 Airflow Connection을 만들어 DAG에 추가해야 해요. Airflow UI나 CLI로 할 수 있어요. 그 일을 해 주는 Airflow 명령 예시예요:
airflow connections add conn_lakefs --conn-type=HTTP --conn-host=http://<LAKEFS_ENDPOINT> \
--conn-extra='{"access_key_id":"<LAKEFS_ACCESS_KEY_ID>","secret_access_key":"<LAKEFS_SECRET_ACCESS_KEY>"}'
lakeFS Airflow 패키지 설치하기
pip로 이 패키지를 설치할 수 있어요
pip install airflow-provider-lakefs
패키지 사용하기
연산자 (Operators)
이 패키지는 lakeFS 서버와 상호작용하는 여러 작업을 노출해요:
CreateBranchOperator는 소스 브랜치(기본값main)에서 새 lakeFS 브랜치를 만들어요.
task_create_branch = CreateBranchOperator(
task_id='create_branch',
repo='example-repo',
branch='example-branch',
source_branch='main'
)
CommitOperator는 커밋되지 않은 변경을 브랜치에 커밋해요.
task_commit = CommitOperator(
task_id='commit',
repo='example-repo',
branch='example-branch',
msg='committing to lakeFS using airflow!',
metadata={'committed_from": "airflow-operator'}
)
MergeOperator는 두 lakeFS 브랜치를 머지해요.
task_merge = MergeOperator(
task_id='merge_branches',
source_ref='example-branch',
destination_branch='main',
msg='merging job outputs',
metadata={'committer': 'airflow-operator'}
)
센서 (Sensors)
실행 중인 DAG와 외부 작업을 동기화할 수 있게 해 주는 센서도 있어요:
CommitSensor는 브랜치에 커밋이 적용될 때까지 기다려요
task_sense_commit = CommitSensor(
repo='example-repo',
branch='example-branch',
task_id='sense_commit'
)
FileSensor는 주어진 파일이 브랜치에 나타날 때까지 기다려요.
task_sense_file = FileSensor(
task_id='sense_file',
repo='example-repo',
branch='example-branch',
path="file/to/sense"
)
예시
airflow-provider-lakeFS 저장소의 이 예시 DAG는 이 모든 것을 사용하는 방법을 보여 줘요.
다른 작업 수행하기
airflow-provider-lakeFS가 아직 지원하지 않는 연산자도 있을 수 있어요. 다음 방법으로 lakeFS에 직접 접근할 수 있어요:
-
SimpleHttpOperator로 lakeFS에 API 요청을 보내요. -
BashOperator와 lakectl 명령을 사용해요.
예를 들어 BashOperator로 브랜치를 삭제하는 예시예요:
commit_extract = BashOperator(
task_id='delete_branch',
bash_command='lakectl branch delete lakefs://example-repo/example-branch',
dag=dag,
)
더 알아보기 (Learn more)
공식 문서: lakeFS Airflow 연동