Apache Spark#

Spark - распределенная вычислительная платформа, предназначенная для обработки больших данных и аналитики. Коннектор Vector DB-Spark позволяет использовать Vector DB в качестве хранилища данных в среде Spark.

Установка#

Для интеграции коннектора в среду Spark получите файл JAR из одного из источников, перечисленных ниже.

  • Выпуски на GitHub

    Файл пакета jar со всеми необходимыми зависимостями можно найти здесь.

  • Сборка из исходного кода

    Чтобы собрать jar из исходников, потребуется установить JDK@8 и Maven. После того как требования будут выполнены, выполните следующую команду в корневой директории проекта.

    mvn package -DskipTests
    

    По умолчанию файл JAR будет записан в директорию target.

  • Центральный репозиторий Maven

    Найдите проект на Maven Central здесь.

Использование#

Создание сеанса Spark с поддержкой Vector DB#

from pyspark.sql import SparkSession

spark = SparkSession.builder.config(
        "spark.jars",
        "path/to/file/spark-VERSION.jar",  # Specify the path to the downloaded JAR file
    )
    .master("local[*]")
    .appName("qdrant")
    .getOrCreate()
import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder
  .config("spark.jars", "path/to/file/spark-VERSION.jar") // Specify the path to the downloaded JAR file
  .master("local[*]")
  .appName("qdrant")
  .getOrCreate()
import org.apache.spark.sql.SparkSession;

public class QdrantSparkJavaExample {
    public static void main(String[] args) {
        SparkSession spark = SparkSession.builder()
                .config("spark.jars", "path/to/file/spark-VERSION.jar") // Specify the path to the downloaded JAR file
                .master("local[*]")
                .appName("qdrant")
                .getOrCreate(); 
    }
}

Загрузка данных#

Перед загрузкой данных с помощью данного коннектора необходимо заранее создать коллекцию с подходящими размерностями вектора и конфигурациями.

Коннектор поддерживает обработку нескольких именованных/неименованных, плотных/разреженных векторов.

Неименованный/вектор по умолчанию#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", <QDRANT_GRPC_URL>)
   .option("collection_name", <QDRANT_COLLECTION_NAME>)
   .option("embedding_field", <EMBEDDING_FIELD_NAME>)  # Expected to be a field of type ArrayType(FloatType)
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Именованный вектор#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", <QDRANT_GRPC_URL>)
   .option("collection_name", <QDRANT_COLLECTION_NAME>)
   .option("embedding_field", <EMBEDDING_FIELD_NAME>)  # Expected to be a field of type ArrayType(FloatType)
   .option("vector_name", <VECTOR_NAME>)
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Примечание

Параметры embedding_field и vector_name поддерживаются для обратной совместимости. Рекомендуется использовать параметры vector_fields и vector_names для именованных векторов, как показано ниже.

Несколько именованных векторов#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("vector_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("vector_names", "<VECTOR_NAME>,<ANOTHER_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Разреженные векторы#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("sparse_vector_value_fields", "<COLUMN_NAME>")
   .option("sparse_vector_index_fields", "<COLUMN_NAME>")
   .option("sparse_vector_names", "<SPARSE_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Несколько разреженных векторов#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("sparse_vector_value_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("sparse_vector_index_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("sparse_vector_names", "<SPARSE_VECTOR_NAME>,<ANOTHER_SPARSE_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Комбинация именованных плотных и разреженных векторов#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("vector_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("vector_names", "<VECTOR_NAME>,<ANOTHER_VECTOR_NAME>")
   .option("sparse_vector_value_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("sparse_vector_index_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("sparse_vector_names", "<SPARSE_VECTOR_NAME>,<ANOTHER_SPARSE_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Мульти-векторы#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("multi_vector_fields", "<COLUMN_NAME>")
   .option("multi_vector_names", "<MULTI_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Несколько мульти-векторов#

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("multi_vector_fields", "<COLUMN_NAME>,<ANOTHER_COLUMN_NAME>")
   .option("multi_vector_names", "<MULTI_VECTOR_NAME>,<ANOTHER_MULTI_VECTOR_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Без векторов#

Вся фреймворк сохраняется как полезная нагрузка:

  <pyspark.sql.DataFrame>
   .write
   .format("io.qdrant.spark.Qdrant")
   .option("qdrant_url", "<QDRANT_GRPC_URL>")
   .option("collection_name", "<QDRANT_COLLECTION_NAME>")
   .option("schema", <pyspark.sql.DataFrame>.schema.json())
   .mode("append")
   .save()

Databricks#

Ознакомьтесь с примером использования коннектора Spark с платформой Databricks.

Можно использовать коннектор qdrant-spark в виде библиотеки в Databricks.

  • Перейдите к разделу Libraries на панели управления кластером Databricks.

  • Выберите пункт Install New, чтобы открыть модальное окно установки библиотек.

  • Введите запрос io.qdrant:spark:VERSION среди пакетов Maven и нажмите кнопку Install.

Databricks

Поддерживаемые типы данных#

Соответствующие типы данных Spark отображаются в полезной нагрузке Vector DB на основе предоставленного параметра schema.

Опции и типы Spark#

Опция

Описание

Тип данных столбца

Обязательно

qdrant_url

URL gRPC экземпляра Vector DB. Например: http://localhost:6334

-

collection_name

Имя коллекции, в которую нужно записать данные

-

schema

Строка JSON схемы фрейма данных

-

embedding_field

Имя столбца, содержащего вложения (устаревший параметр - используйте вместо него vector_fields)

ArrayType(FloatType)

id_field

Имя столбца, содержащего идентификаторы точек. По умолчанию: случайный UUID

StringType или IntegerType

batch_size

Максимальный размер партии загрузки. По умолчанию: 64

-

retries

Количество повторных попыток загрузки. По умолчанию: 3

-

api_key

Ключ API Vector DB для аутентификации

-

vector_name

Имя вектора в коллекции

-

vector_fields

Запятая-разделенный список названий столбцов, содержащих векторы

ArrayType(FloatType)

vector_names

Запятая-разделенный список названий векторов в коллекции

-

sparse_vector_index_fields

Запятая-разделенный список названий столбцов, содержащих индексы разреженных векторов

ArrayType(IntegerType)

sparse_vector_value_fields

Запятая-разделенный список названий столбцов, содержащих значения разреженных векторов

ArrayType(FloatType)

sparse_vector_names

Запятая-разделенный список названий разреженных векторов в коллекции

-

multi_vector_fields

Запятая-разделенный список названий столбцов, содержащих значения многомерных векторов

ArrayType(ArrayType(FloatType))

multi_vector_names

Запятая-разделенный список названий многомерных векторов в коллекции

-

shard_key_selector

Запятая-разделенный список пользовательских ключей шардов, используемых при вставке

-

wait

Ожидать завершения каждой партии вставки. true или false. Значение по умолчанию: true

-

Для получения дополнительной информации обязательно ознакомьтесь с репозиторием Qdrant-Spark на GitHub. Руководство по Apache Spark доступно в официальной документации.