Настройка потоковой передачи данных с использованием Kafka через Confluent#

Автор: М.К.Паван Кумар, исследователь из ИИИТДМ, Курнул. Специалист по методам смягчения галлюцинаций и методологиям RAG. • GitHubMedium

Введение#

В данном руководстве пошагово описан через процесс установки и настройки коннектора приемника Platform V Vector DB (далее - Vector DB), создания необходимой инфраструктуры и разработки практического демонстрационного приложения. Изучив документ сформируется глубокое понимание того, как использовать эту мощную интеграцию для оптимизации рабочих процессов с данными, тем самым повышая производительность и возможности приложений реального времени семантического поиска и RAG, основанных на данных.

В этом примере исходные данные будут поступать из хранилища BLOB-объектов Azure и MongoDB.

Реальное время захвата изменений данных (CDC) с Kafka и Platform V Vector DB (далее - Vector DB):

1.webp

Архитектура#

Источники данных#

Архитектура начинается с источников данных, представленных MongoDB и хранилищем BLOB-объектов Azure. Эти системы играют ключевую роль в хранении и управлении сырыми данными. MongoDB, популярная база данных NoSQL, известна гибкостью при работе с различными форматами данных и способностью к горизонтальному масштабированию. Она широко используется в приложениях, требующих высокой производительности и масштабируемости. Хранилище BLOB-объектов Azure, с другой стороны, является решением от Microsoft для облачного хранения объектов. Оно предназначено для хранения огромных объемов неструктурированных данных, таких как текстовые или двоичные данные. Данные извлекаются из этих источников с помощью соединителей источника, отвечающих за захват изменений в реальном времени и их отправку в Kafka.

Kafka#

В центре этой архитектуры находится Kafka, распределенная платформа потоковой обработки событий, способная обрабатывать триллионы событий в день. Kafka выступает в роли центрального хаба, где данные из различных источников могут быть поглощены, обработаны и переданы различным системам нижнего уровня. Ее отказоустойчивая и масштабируемая архитектура гарантирует надежную передачу и обработку данных в режиме реального времени. Способность Kafka справляться с высокопроизводительными низкоотложенными потоками данных делает ее идеальным выбором для обработки и аналитики данных в реальном времени. Использование Confluent расширяет функциональные возможности Kafka, предоставляя дополнительные инструменты и сервисы для управления кластерами Kafka и потоковой обработки.

Vector DB#

Обработанные данные затем направляются в Vector DB, высокоэффективный векторный поисковый движок, предназначенный для поисковых запросов по подобию. Vector DB отлично справляется с управлением и поиском многомерных векторных данных, что имеет решающее значение для приложений, связанных с машинным обучением и искусственным интеллектом, такими как рекомендательные системы, распознавание изображений и обработка естественного языка. Коннектор-приемник Vector DB для Kafka играет здесь ключевую роль, обеспечивая плавную интеграцию между Kafka и Vector DB. Этот коннектор позволяет осуществлять прием векторных данных в Vector DB в реальном времени, гарантируя актуальность данных и готовность к высокопроизводительным запросам по сходству.

Важность интеграции и конвейера#

Интеграция этих компонентов формирует мощный и эффективный конвейер потоковой передачи данных. Коннектор-приемник Vector DB обеспечивает непрерывное поступление данных, проходящих через Kafka, в Vector DB без какого-либо ручного вмешательства. Эта интеграция в реальном времени критически важна для приложений, зависящих от самых свежих данных для принятия решений и анализа. Комбинируя сильные стороны MongoDB и хранилища BLOB-объектов Azure для хранения данных, Kafka для потоковой передачи данных и Vector DB для векторного поиска, этот конвейер предоставляет надежное решение для управления и обработки больших объемов данных в реальном времени. Масштабируемость, устойчивость к сбоям и способность к обработке данных в реальном времени являются ключевыми факторами эффективности данной архитектуры, делая ее универсальным решением для современных приложений, ориентированных на работу с данными.

Установка платформы Confluent Kafka#

Чтобы установить платформу Confluent Kafka (самостоятельно управляемую локально), выполните следующие три простых шага:

Загрузите и распакуйте дистрибутив:

  • Посетите страницу установки Confluent.

  • Загрузите файлы дистрибутива (tar, zip и др.).

  • Распакуйте загруженный файл с помощью команды:

    tar -xvf confluent-<version>.tar.gz
    

    или

    unzip confluent-<version>.zip
    
  • Настройте переменные среды:

    # Set CONFLUENT_HOME to the installation directory:
    export CONFLUENT_HOME=/path/to/confluent-<version>
    
    # Add Confluent binaries to your PATH
    export PATH=$CONFLUENT_HOME/bin:$PATH
    
  • Запустите платформу Confluent локально:

    # Start the Confluent Platform services:
    confluent local start
    # Stop the Confluent Platform services:
    confluent local stop
    

Установка Vector DB#

Для установки и запуска Vector DB (локально самостоятельно управляемого) обратитесь к руководству по установке.

Установка коннектора-приемника Vector DB-Kafka#

Для установки коннектора Vector DB Kafka с помощью Confluent Hub можно использовать простую команду confluent-hub install. Эта команда упрощает процесс, устраняя необходимость вручную изменять конфигурационные файлы. Чтобы установить версию коннектора Vector DB Kafka 1.1.0, выполните следующую команду в терминале:

 confluent-hub install qdrant/qdrant-kafka:1.1.0

Эта команда скачивает и устанавливает указанный коннектор непосредственно из Confluent Hub в среду Confluent Platform или Kafka Connect. Процесс установки автоматически обрабатывает все необходимые зависимости, обеспечивая гладкую интеграцию коннектора Vector DB Kafka с существующей настройкой. После установки коннектором можно управлять и настраивать с помощью Центра управления Confluent или REST API Kafka Connect, позволяя эффективно передавать данные между Kafka и Vector DB без необходимости сложной ручной настройки.

Локальная платформа Confluent после установки показывает соединители источника и приемника:

2.webp

Убедитесь, что после установки правильно настроен коннектор следующим образом. Обратите внимание, что параметры key.converter и value.converter очень важны для безопасной доставки сообщений Kafka из темы в Vector DB.

{
  "name": "QdrantSinkConnectorConnector_0",
  "config": {
    "value.converter.schemas.enable": "false",
    "name": "QdrantSinkConnectorConnector_0",
    "connector.class": "io.qdrant.kafka.QdrantSinkConnector",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "topics": "topic_62,qdrant_kafka.docs",
    "errors.deadletterqueue.topic.name": "dead_queue",
    "errors.deadletterqueue.topic.replication.factor": "1",
    "qdrant.grpc.url": "http://localhost:6334",
    "qdrant.api.key": "************"
  }
}

Установка MongoDB#

Для подключения Kafka к MongoDB в качестве источника экземпляр MongoDB должен работать в режиме replicaSet. Ниже приведен файл конфигурации docker compose, который запускает одиночный узел экземпляра MongoDB в режиме replicaSet.

version: "3.8"

services:
  mongo1:
    image: mongo:7.0
    command: ["--replSet", "rs0", "--bind_ip_all", "--port", "27017"]
    ports:
      - 27017:27017
    healthcheck:
      test: echo "try { rs.status() } catch (err) { rs.initiate({_id:'rs0',members:[{_id:0,host:'host.docker.internal:27017'}]}) }" | mongosh --port 27017 --quiet
      interval: 5s
      timeout: 30s
      start_period: 0s
      start_interval: 1s
      retries: 30
    volumes:
      - "mongo1_data:/data/db"
      - "mongo1_config:/data/configdb"

volumes:
  mongo1_data:
  mongo1_config:

Аналогичным образом установите и сконфигурируйте источник-коннектор следующим образом.

confluent-hub install mongodb/kafka-connect-mongodb:latest

После установки коннектора MongoDB конфигурация коннектора должна выглядеть так:

{
  "name": "MongoSourceConnectorConnector_0",
  "config": {
    "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "connection.uri": "mongodb://127.0.0.1:27017/?replicaSet=rs0&directConnection=true",
    "database": "qdrant_kafka",
    "collection": "docs",
    "publish.full.document.only": "true",
    "topic.namespace.map": "{\"*\":\"qdrant_kafka.docs\"}",
    "copy.existing": "true"
  }
}

Демонстрационное приложение#

Теперь, когда инфраструктура полностью готова, пришло время создать простое приложение и проверить установку. Цель приложения заключается в том, чтобы вставить данные в MongoDB, а затем они были бы также приняты в Vector DB с использованием захвата изменений данных (CDC).

requirements.txt

fastembed==0.3.1
pymongo==4.8.0
qdrant_client==1.10.1

project_root_folder/main.py

Это всего лишь пример кода. Тем не менее он может быть расширен до миллионов операций в соответствии со сценарием использования.

from pymongo import MongoClient
from utils.app_utils import create_qdrant_collection
from fastembed import TextEmbedding

collection_name: str = 'test'
embed_model_name: str = 'snowflake/snowflake-arctic-embed-s'
# Step 0: create qdrant_collection
create_qdrant_collection(collection_name=collection_name, embed_model=embed_model_name)

# Step 1: Connect to MongoDB
client = MongoClient('mongodb://127.0.0.1:27017/?replicaSet=rs0&directConnection=true')

# Step 2: Select Database
db = client['qdrant_kafka']

# Step 3: Select Collection
collection = db['docs']

# Step 4: Create a Document to Insert

description = "qdrant is a high available vector search engine"
embedding_model = TextEmbedding(model_name=embed_model_name)
vector = next(embedding_model.embed(documents=description)).tolist()
document = {
    "collection_name": collection_name,
    "id": 1,
    "vector": vector,
    "payload": {
        "name": "qdrant",
        "description": description,
        "url": "https://qdrant.tech/documentation"
    }
}

# Step 5: Insert the Document into the Collection
result = collection.insert_one(document)

# Step 6: Print the Inserted Document ID
print("Inserted document ID:", result.inserted_id)

project_root_folder/utils/app_utils.py

from qdrant_client import QdrantClient, models

client = QdrantClient(url="http://localhost:6333", api_key="<YOUR_KEY>")
dimension_dict = {"snowflake/snowflake-arctic-embed-s": 384}

def create_qdrant_collection(collection_name: str, embed_model: str):

    if not client.collection_exists(collection_name=collection_name):
        client.create_collection(
            collection_name=collection_name,
            vectors_config=models.VectorParams(size=dimension_dict.get(embed_model), distance=models.Distance.COSINE)
        )

Перед запуском приложения ниже представлено состояние баз данных MongoDB и Vector DB.

Исходное состояние: нет коллекции под названием test & no data в коллекции docs базы данных MongoDB:

3.webp

Когда запускается код, данные поступают в MongoDB, происходит срабатывание CDC, и в конечном итоге Vector DB получает эти данные.

Тестовая коллекция Vector DB создается автоматически:

4.webp

Данные вставлены как в MongoDB, так и в Vector DB:

5.webp

Заключение#

Подводя итог, интеграция Kafka с Vector DB с использованием коннектора приемника Vector DB предоставляет бесшовное и эффективное решение для потоковой передачи и обработки данных в реальном времени. Такая настройка не только повышает возможности конвейера данных, но и гарантирует, что многомерные векторные данные постоянно индексируются и доступны для поиска по сходству. Следуя руководству по установке и настройке, можно легко организовать надежный поток данных из источников данных, таких как MongoDB и хранилище BLOB-объектов Azure, через Kafka и в Vector DB. Данная архитектура дает современным приложениям возможность использовать аналитику данных в реальном времени и передовые функции поиска, прокладывая путь к инновационным решениям, основанным на данных.