Apache Airflow 연동

Apache Airflow 연동 (airflow)

Apache Airflow는 데이터·컴퓨팅 워크플로우를 작성하고 스케줄링하며 모니터링하기 위한 오픈소스 플랫폼이에요. Airflow는 Python을 사용해서 쉽게 스케줄링하고 모니터링할 수 있는 워크플로우를 만들죠. 이번에는 Airflow에서 Qdrant를 연동하는 방법을 함께 살펴볼게요.

Airflow에서는 Qdrant를 provider 형태로 제공하고 있어요. 이 provider를 사용하면 파이프라인 안에서 Qdrant 데이터베이스를 직접 다룰 수 있습니다.

출처: Qdrant 공식 문서 — airflow

사전 준비 (Prerequisites)

Airflow를 설정하기 전에 필요한 것들이 두 가지 있어요.

  1. 연결할 Qdrant 인스턴스 — 설치 가이드에서 설정할 수 있어요.
  2. 실행 중인 Airflow 인스턴스 — Airflow Quick Start Guide를 참고하면 돼요.

설치 (Installation)

Airflow 셸에서 다음 명령으로 Qdrant provider를 설치할 수 있어요.

pip install apache-airflow-providers-qdrant

참고: provider를 사용할 수 있으려면 Airflow 세션을 재시작해야 해요.

연결(Connection) 설정하기

Airflow UI의 Admin -> Connections 섹션을 열고, Create 링크를 클릭해서 새 Qdrant connection을 만들어요.

연결은 환경 변수외부 시크릿 백엔드를 통해서도 설정할 수 있어요.

Qdrant Hook

Airflow의 hook은 특정 API를 추상화한 것으로, Airflow가 외부 시스템과 상호작용할 수 있게 해 주는 개념이에요. Qdrant hook도 같은 역할을 하고요.

from airflow.providers.qdrant.hooks.qdrant import QdrantHook

hook = QdrantHook(conn_id="qdrant_connection")
hook.verify_connection()

QdrantHook 인스턴스의 @property conn을 통해서 qdrant_client#QdrantClient 인스턴스에 접근할 수 있어요. 이 인스턴스를 Airflow 워크플로우 안에서 바로 사용하면 됩니다.

from qdrant_client import models

hook.conn.count("<COLLECTION_NAME>")
hook.conn.upsert(
    "<COLLECTION_NAME>",
    points=[
        models.PointStruct(
            id=32,
            vector=[0.32, 0.12, 0.123],
            payload={"color": "red"},
        ),
    ],
)

Qdrant Ingest Operator

Qdrant provider에는 Qdrant 컬렉션에 데이터를 업로드하는 편리한 operator도 포함되어 있어요. 내부적으로 앞서 본 Qdrant hook을 사용하죠.

from airflow.providers.qdrant.operators.qdrant import QdrantIngestOperator

vectors = [
    [0.11, 0.22, 0.33, 0.44],
    [0.55, 0.66, 0.77, 0.88],
    [0.88, 0.11, 0.12, 0.13],
]
ids = [32, 21, "b626f6a9-b14d-4af9-b7c3-43d8deb719a6"]
payload = [
    {"meta": "data"},
    {"meta": "data_2"},
    {"meta": "data_3", "extra": "data"},
]

QdrantIngestOperator(
    conn_id="qdrant_connection",
    task_id="qdrant_ingest",
    collection_name="<COLLECTION_NAME>",
    vectors=vectors,
    ids=ids,
    payload=payload,
)

참고 자료 (Reference)

더 알아보기 (Learn more)