Представляем Databricks Connector — прозрачное решение для оптимизации миграции и преобразования неструктурированных данных
В условиях стремительного развития технологий ИИ и машинного обучения (AI/ML) векторные эмбеддинги стали предпочтительным методом для индексирования неструктурированных данных и выполнения семантического поиска.
Поиск на базе ИИ обычно состоит из двух отдельных этапов: офлайн-индексирования данных и онлайн-обслуживания запросов, что требует использования разных технологических стеков. Однако попытка эффективно и бесшовно перенести данные из стека офлайн-обработки в стек онлайн-обслуживания может быть сложной задачей. Последний релиз Zilliz Cloud представляет Databricks Connector, продуманное решение для оптимизации этого процесса за счет интеграции Apache Spark/Databricks и Milvus/Zilliz Cloud.
В этом посте мы представим эту интеграцию, рассмотрим ее практическое применение в реальных сценариях и пошагово покажем, как ею пользоваться.
Как работает Databricks Connector и его сценарии использования
Spark известен своей способностью обрабатывать крупномасштабные данные и эффективностью в машинном обучении. С другой стороны, Milvus отлично справляется с эффективной обработкой и поиском векторных эмбеддингов, генерируемых моделями машинного обучения. Сочетание этих двух мощных технологий способствует разработке передовых приложений в таких областях, как генеративный ИИ, рекомендательные системы и поиск изображений и видео.
Передача данных из Spark в Milvus — распространенная задача при создании поиска на базе ИИ, но она часто требует сложного связующего кода в бэкенде поисковой системы. Коннектор Spark-Milvus упрощает этот процесс, сводя его к одному вызову функции внутри программы Spark.
Этот коннектор полезен в различных сценариях.
Пакетный импорт данных
Команды с экспертизой в машинном обучении часто обновляют свои модели эмбеддингов, чтобы учитывать новейшие результаты исследований. После каждого обновления модели эмбеддингов, требующего повторной обработки всего корпуса данных заданием Spark и генерации нового набора векторов, командам часто нужен пользовательский «связующий код», выделенный сервис или другое задание Spark для интеграции этих новых векторов в стек обслуживания. Однако с Databricks Connector эта задача становится такой же простой, как предоставление заданию доступа на запись в S3 bucket Milvus (или во временный bucket при использовании Zilliz Cloud). Этот оптимизированный процесс позволяет заданию Spark, генерирующему векторы, загружать данные напрямую в экземпляр Milvus через простой вызов утилитарной функции.
Итеративная вставка
Пользователи, работающие не в больших масштабах, также могут извлечь пользу из прямой вставки записей Spark DataFrame в Milvus с помощью коннектора Spark-Milvus. Такой подход экономит усилия на написание кода для установления соединения и вызовов API, делая процесс интеграции более плавным.
Как использовать Databricks Connector
В этом разделе демонстрируется использование Databricks Connector для оптимизации миграции и преобразования данных.
Итеративная вставка Spark Dataframe
Благодаря подключению Spark-Milvus потоковая передача данных из Spark в Milvus стала проще, чем когда-либо. Теперь вы можете напрямую отправлять данные из Spark в Milvus с помощью нативного Dataframe API Spark. Тот же код также работает для Databricks и Zilliz (полностью управляемого Milvus). Вот фрагмент кода, демонстрирующий этот подход:
// Specify the target Milvus instance and vector data collection
df.write.format("milvus")
.option(MILVUS_URI, "https://in01-xxxxxxxxx.aws-us-west-2.vectordb.zillizcloud.com:19535")
.option(MILVUS_TOKEN, dbutils.secrets.get(scope = "zillizcloud", key = "token"))
.option(MILVUS_COLLECTION_NAME, "text_embedding")
.option(MILVUS_COLLECTION_VECTOR_FIELD, "embedding")
.option(MILVUS_COLLECTION_VECTOR_DIM, "128")
.option(MILVUS_COLLECTION_PRIMARY_KEY, "id")
.mode(SaveMode.Append)
.save()
Пакетная загрузка коллекции
Мы рекомендуем использовать функцию `MilvusUtils. bulkInsertFromSpark ()` в ситуациях, когда необходимо эффективно передавать большие объемы данных. Этот подход чрезвычайно эффективен для обработки огромных наборов данных.
Подход для Milvus
Интеграция включает в себя бакеты S3 или MinIO для самостоятельно размещенных экземпляров Milvus в качестве внутреннего хранилища. Предоставив доступ Spark или Databricks, задание Spark может использовать коннекторы Milvus для пакетной записи данных в бакет, а затем выполнить массовую вставку всей коллекции для обслуживания.
// Write the data in batch into the Milvus bucket storage.
val outputPath = "s3a://milvus-bucket/result"
df.write
.mode("overwrite")
.format("parquet")
.save(outputPath)
// Specify Milvus options.
val targetProperties = Map(
MilvusOptions.MILVUS_HOST -> host,
MilvusOptions.MILVUS_PORT -> port.toString,
MilvusOptions.MILVUS_COLLECTION_NAME -> targetCollectionName,
MilvusOptions.MILVUS_BUCKET -> bucketName,
MilvusOptions.MILVUS_ROOTPATH -> rootPath,
MilvusOptions.MILVUS_FS -> fs,
MilvusOptions.MILVUS_STORAGE_ENDPOINT -> minioEndpoint,
MilvusOptions.MILVUS_STORAGE_USER -> minioAK,
MilvusOptions.MILVUS_STORAGE_PASSWORD -> minioSK,
)
val targetMilvusOptions = new MilvusOptions(new CaseInsensitiveStringMap(targetProperties.asJava))
// Bulk insert Spark output files into Milvus
MilvusUtils.bulkInsertFromSpark(spark, targetMilvusOptions, outputPath, "parquet")
Подход для Zilliz Cloud
Если вы используете Zilliz Cloud (управляемый Milvus), вы можете воспользоваться его удобным Data Import API. Zilliz Cloud предоставляет комплексные инструменты и документацию, которые помогут вам эффективно переносить данные из различных источников данных, включая Spark. Настроив бакет S3 в качестве промежуточного хранилища и предоставив доступ Zilliz Cloud, Data Import API без проблем загружает данные из бакета S3 в векторную базу данных.
Перед запуском этой интеграции необходимо загрузить среду выполнения Spark, добавив jar-файл в Databricks Cluster. Существуют разные способы установки библиотеки. На скриншоте ниже показана загрузка jar-файла с локального компьютера в кластер.
Для получения более подробной информации об установке библиотеки в рабочей области Databricks обратитесь к официальной документации Databricks, чтобы узнать больше.
Массовая вставка требует хранения данных во временном бакете, чтобы Zilliz Cloud мог импортировать их пакетами. Вы можете создать бакет S3 и настроить его как внешнее расположение Databricks. Подробности см. в этой документации. Чтобы снизить риск безопасности для учетных данных Zilliz Cloud, вы можете безопасно управлять ими в Databricks, следуя инструкциям Databricks.
Вот фрагмент кода, демонстрирующий процесс пакетной миграции данных. Как и в приведенном выше примере Milvus, вам нужно только заменить учетные данные и адрес бакета S3.
// Запишите данные пакетно в хранилище bucket Milvus.
val outputPath = "s3://my-temp-bucket/result"
df.write
.mode("overwrite")
.format("mjson")
.save(outputPath)
// Укажите параметры Milvus.
val targetProperties = Map(
MilvusOptions.MILVUS_URI -> zilliz_uri,
MilvusOptions.MILVUS_TOKEN -> zilliz_token,
MilvusOptions.MILVUS_COLLECTION_NAME -> targetCollectionName,
MilvusOptions.MILVUS_BUCKET -> bucketName,
MilvusOptions.MILVUS_ROOTPATH -> rootPath,
MilvusOptions.MILVUS_FS -> fs,
MilvusOptions.MILVUS_STORAGE_ENDPOINT -> minioEndpoint,
MilvusOptions.MILVUS_STORAGE_USER -> minioAK,
MilvusOptions.MILVUS_STORAGE_PASSWORD -> minioSK,
)
val targetMilvusOptions = new MilvusOptions(new CaseInsensitiveStringMap(targetProperties.asJava))
// Массовая вставка выходных файлов Spark в Milvus
MilvusUtils.bulkInsertFromSpark(spark, targetMilvusOptions, outputPath, "mjson")
Собираем всё вместе: пример notebook, который проведет вас через весь процесс
Чтобы помочь вам быстро начать, мы подготовили пример notebook, который проведет вас через процессы потоковой и пакетной передачи данных с Milvus и Zilliz Cloud.
Заключение
Интеграция Spark и Milvus открывает захватывающие возможности для приложений на базе ИИ. Благодаря нашему упрощенному подходу к переносимости данных разработчики могут без труда передавать данные из Spark/Databricks в Milvus/Zilliz Cloud — как в реальном времени, так и в пакетном режиме. Эта интеграция позволяет создавать эффективные и масштабируемые ИИ-решения, раскрывая весь потенциал этих мощных технологий.
Готовы отправиться в свое ИИ-путешествие? Начните бесплатно с Zilliz Cloud уже сегодня — без сложностей с установкой и без необходимости указывать кредитную карту.
Читать далее

What Is a Vector Lakebase?
A Vector Lakebase is a unified, lake-native data architecture for AI that combines vector-database-grade serving with open lake storage, reusable lake-level indexes, and a shared semantic layer.

VidTok: Rethinking Video Processing with Compact Tokenization
VidTok tokenizes videos to reduce redundancy while preserving spatial and temporal details for efficient processing.

Vector Databases vs. Key-Value Databases
Use a vector database for AI-powered similarity search; use a key-value database for high-throughput, low-latency simple data lookups.



