Vector DB на платформе Databricks#
Databricks - унифицированная аналитическая платформа для работы с большими данными и искусственным интеллектом. Она построена вокруг Apache Spark, мощной открытой распределенной вычислительной системы, хорошо подходящей для обработки крупномасштабных наборов данных и выполнения сложных аналитических задач.
Apache Spark спроектирована так, чтобы масштабироваться горизонтально, то есть она может обрабатывать дорогостоящие операции, такие как генерация векторных вложений, распределяя вычисления по кластеру машин. Эта масштабируемость критически важна при работе с большими наборами данных.
В этом примере представлено, как векторизовать набор данных с плотными и разреженными вложениями с помощью библиотеки FastEmbed от Platform V Vector DB (далее - Vector DB). Затем загрузитн этот векторизованный набор данных в кластер Vector DB с использованием коннектора Vector DB Spark на платформе Databricks.
Настройка проекта на платформе Databricks#
Настройте кластер Databricks, следуя официальным рекомендациям документации.
Установите библиотеку коннектор Vector DB Spark следующим образом:
Перейдите в раздел
Librariesна панели управления кластером.Нажмите кнопку
Install Newв правом верхнем углу, чтобы открыть модальное окно установки библиотеки.Найдите пакет
io.qdrant:spark:VERSIONв репозитории Maven и нажмитеInstall.

Создайте новую тетрадь Databricks на кластере, чтобы начать работу с данными и библиотеками.
Загрузка набора данных#
Установите необходимые зависимости:
%pip install fastembed datasetsСкачайте набор данных:
from datasets import load_dataset dataset_name = "tasksource/med" dataset = load_dataset(dataset_name, split="train") # We'll use the first 100 entries from this dataset and exclude some unused columns. dataset = dataset.select(range(100)).remove_columns(["gold_label", "genre"])Преобразуйте набор данных в датафрейм Spark:
dataset.to_parquet("/dbfs/pq.pq") dataset_df = spark.read.parquet("file:/dbfs/pq.pq")
Векторизация данных#
В данном разделе будут сгенерированы как плотные, так и разреженные векторы для записей с помощью FastEmbed. Будет создана пользовательская функцию (UDF), которая будет выполнять эту задачу.
Создание функции векторизации#
from fastembed import TextEmbedding, SparseTextEmbedding
def vectorize(partition_data):
# Initialize dense and sparse models
dense_model = TextEmbedding(model_name="BAAI/bge-small-en-v1.5")
sparse_model = SparseTextEmbedding(model_name="Qdrant/bm25")
for row in partition_data:
# Generate dense and sparse vectors
dense_vector = next(dense_model.embed(row.sentence1))
sparse_vector = next(sparse_model.embed(row.sentence2))
yield [
row.sentence1, # 1st column: original text
row.sentence2, # 2nd column: original text
dense_vector.tolist(), # 3rd column: dense vector
sparse_vector.indices.tolist(), # 4th column: sparse vector indices
sparse_vector.values.tolist(), # 5th column: sparse vector values
]
Используется модель BAAI/bge-small-en-v1.5 для плотных вложений и метод BM25 для разреженных вложений.
Применение UDF к датафрейму#
Далее примените vectorize UDF к датафрейму Spark для создания вложений.
embeddings = dataset_df.rdd.mapPartitions(vectorize)
Метод mapPartitions() возвращает устойчивый распределенный набор данных (RDD), который затем должен быть преобразован обратно в датафрейм Spark.
Создание нового датафрейма Spark с векторизованными данными#
Создайте новый датафрейм Spark (embeddings_df) с векторизованными данными, используя заданную схему.
from pyspark.sql.types import StructType, StructField, StringType, ArrayType, FloatType, IntegerType
# Define the schema for the new dataframe
schema = StructType([
StructField("sentence1", StringType()),
StructField("sentence2", StringType()),
StructField("dense_vector", ArrayType(FloatType())),
StructField("sparse_vector_indices", ArrayType(IntegerType())),
StructField("sparse_vector_values", ArrayType(FloatType()))
])
# Create the new dataframe with the vectorized data
embeddings_df = spark.createDataFrame(data=embeddings, schema=schema)
Загрузка данных в Vector DB#
Создайте коллекцию Vector DB:
Следуйте руководству по созданию коллекции с соответствующими конфигурациями. Вот пример запроса для поддержки как плотных, так и разреженных векторов:
PUT /collections/{collection_name} { "vectors": { "dense": { "size": 384, "distance": "Cosine" } }, "sparse_vectors": { "sparse": {} } }
Загрузите датафрейм в Qdrant:
options = { "qdrant_url": "<QDRANT_GRPC_URL>", "api_key": "<QDRANT_API_KEY>", "collection_name": "<QDRANT_COLLECTION_NAME>", "vector_fields": "dense_vector", "vector_names": "dense", "sparse_vector_value_fields": "sparse_vector_values", "sparse_vector_index_fields": "sparse_vector_indices", "sparse_vector_names": "sparse", "schema": embeddings_df.schema.json(), } embeddings_df.write.format("io.qdrant.spark.Qdrant").options(**options).mode( "append" ).save()
Не забудьте заменить значения-заполнители (<QDRANT_GRPC_URL>, <QDRANT_API_KEY>, <QDRANT_COLLECTION_NAME>) реальными значениями. Если параметр id_field не указан, коннектор Vector DB Spark генерирует случайные идентификаторы UUID для каждой точки.
Результат команды должен выглядеть так:
Command took 40.37 seconds -- by xxxxx90@xxxxxx.com at 4/17/2024, 12:13:28 PM on fastembed