Apache Airflow#

Apache Airflow - открытая платформа для создания, планирования и мониторинга рабочих процессов обработки данных и вычислений. Airflow использует Python для создания рабочих процессов, которые можно легко планировать и контролировать.

Vector DB доступен в качестве провайдера в Airflow для взаимодействия с базой данных.

Предварительные требования#

Перед настройкой Airflow необходимо иметь:

  1. Экземпляр Vector DB для подключения. Можно настроить его в руководстве по установке.

  2. Запущенный экземпляр Airflow. Можно воспользоваться их Быстрым стартом.

Установка#

Можно установить провайдер Vector DB, выполнив команду pip install apache-airflow-providers-qdrant в оболочке Airflow.

Примечание

Потребуется перезапустить сеанс Airflow, чтобы провайдер стал доступным.

Настройка соединения#

Откройте раздел Admin-> Connections интерфейса пользователя Airflow. Нажмите ссылку Create, чтобы создать новое соединение с Vector DB.

Соединение с

Вы также можете настроить соединение с помощью переменных окружения или внешнего хранилища секретов .

Хук Vector DB#

Хук Airflow представляет собой абстракцию конкретного API, позволяющую Airflow взаимодействовать с внешней системой.

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

hook = QdrantHook(conn_id="qdrant_connection")

hook.verify_connection()

Экземпляр qdrant_client#QdrantClient доступен через @property conn экземпляра QdrantHook для использования внутри рабочих процессов 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"})
    ],
)

Оператор загрузки Vector DB#

Провайдер Vector DB также предоставляет удобный оператор для выгрузки данных в коллекцию Vector DB, который внутренне использует хук Vector DB.

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,
)

Справочные материалы#