Настройка потоковой передачи данных с использованием Kafka через Confluent#
Автор: М.К.Паван Кумар, исследователь из ИИИТДМ, Курнул. Специалист по методам смягчения галлюцинаций и методологиям RAG. • GitHub • Medium
Введение#
В данном руководстве пошагово описан через процесс установки и настройки коннектора приемника Platform V Vector DB (далее - Vector DB), создания необходимой инфраструктуры и разработки практического демонстрационного приложения. Изучив документ сформируется глубокое понимание того, как использовать эту мощную интеграцию для оптимизации рабочих процессов с данными, тем самым повышая производительность и возможности приложений реального времени семантического поиска и RAG, основанных на данных.
В этом примере исходные данные будут поступать из хранилища BLOB-объектов Azure и MongoDB.
Реальное время захвата изменений данных (CDC) с Kafka и Platform V Vector DB (далее - Vector DB):

Архитектура#
Источники данных#
Архитектура начинается с источников данных, представленных 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 после установки показывает соединители источника и приемника:

Убедитесь, что после установки правильно настроен коннектор следующим образом. Обратите внимание, что параметры 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:

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

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

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