Apache Spark 연동

Apache Spark 연동 (spark)

Spark는 빅데이터 처리와 분석을 위해 설계된 분산 컴퓨팅 프레임워크예요. Qdrant-Spark 커넥터를 사용하면 Spark에서 Qdrant를 storage destination(저장 목적지)으로 사용할 수 있게 됩니다.

출처: Qdrant 공식 문서 — spark

설치 (Installation)

커넥터를 Spark 환경에 통합하려면, 아래 나열된 소스 중 하나에서 JAR 파일을 받으면 돼요.

GitHub Releases

모든 필수 의존성이 포함된 패키징된 jar 파일은 여기에서 찾을 수 있어요.

소스에서 빌드하기

jar을 소스에서 빌드하려면 JDK@8Maven이 설치되어 있어야 해요. 요구사항이 충족되면 프로젝트 루트에서 다음 명령을 실행하세요.

mvn package -DskipTests

JAR 파일은 기본적으로 target 디렉토리에 작성돼요.

Maven Central

Maven Central에서 프로젝트를 찾으려면 여기를 확인하세요.

사용법 (Usage)

Qdrant 지원이 포함된 Spark 세션 만들기

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .config("spark.jars", "path/to/file/spark-VERSION.jar")  # 다운로드한 JAR 파일 경로 지정
    .master("local[*]")
    .appName("qdrant")
    .getOrCreate()
)
import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder
  .config("spark.jars", "path/to/file/spark-VERSION.jar") // 다운로드한 JAR 파일 경로 지정
  .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") // 다운로드한 JAR 파일 경로 지정
            .master("local[*]")
            .appName("qdrant")
            .getOrCreate();
    }
}

데이터 로드 (Loading data)

이 커넥터를 사용해 데이터를 로드하기 전에, 적절한 벡터 차원과 설정을 가진 컬렉션이 미리 생성되어 있어야 해요.

커넥터는 여러 개의 named/unnamed, dense/sparse 벡터를 수집(ingest)할 수 있어요.

Unnamed/Default 벡터

<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>) \  # ArrayType(FloatType) 타입 필드여야 함
    .option("schema", <pyspark.sql.DataFrame>.schema.json()) \
    .mode("append") \
    .save()

Named 벡터

<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>) \  # ArrayType(FloatType) 타입 필드여야 함
    .option("vector_name", <VECTOR_NAME>) \
    .option("schema", <pyspark.sql.DataFrame>.schema.json()) \
    .mode("append") \
    .save()

참고

embedding_fieldvector_name 옵션은 하위 호환성을 위해 유지되고 있어요. named 벡터의 경우 아래와 같이 vector_fieldsvector_names를 사용하는 것이 권장됩니다.

여러 개의 named 벡터

<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()

Sparse 벡터

<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()

여러 개의 sparse 벡터

<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()

Named dense 벡터와 sparse 벡터의 조합

<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()

Multi-vectors

<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()

여러 개의 Multi-vectors

<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()

벡터 없이 — 전체 데이터프레임을 payload로 저장

<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

Databricks에서 qdrant-spark 커넥터를 라이브러리로 사용할 수 있어요.

  • Databricks 클러스터 대시보드의 Libraries 섹션으로 이동해요.
  • Install New를 선택해 라이브러리 설치 모달을 열어요.
  • Maven 패키지에서 io.qdrant:spark:VERSION을 검색하고 Install을 클릭하세요.

Datatype 지원

제공된 schema에 따라 적절한 Spark 데이터 타입이 Qdrant payload에 매핑돼요.

옵션과 Spark 타입

Option 설명 Column DataType Required
qdrant_url Qdrant 인스턴스의 gRPC URL. 예: http://localhost:6334 -
collection_name 데이터를 쓸 컬렉션 이름 -
schema 데이터프레임 스키마의 JSON 문자열 -
embedding_field 임베딩을 담은 컬럼 이름 (Deprecated — vector_fields 사용) ArrayType(FloatType)
id_field Point ID를 담은 컬럼 이름. 기본값: Random UUID StringType 또는 IntegerType
batch_size 업로드 배치의 최대 크기. 기본값: 64 -
retries 업로드 재시도 횟수. 기본값: 3 -
api_key 인증용 Qdrant API 키 -
vector_name 컬렉션 내 벡터 이름. -
vector_fields 벡터를 담은 컬럼 이름들(콤마 구분). ArrayType(FloatType)
vector_names 컬렉션 내 벡터 이름들(콤마 구분). -
sparse_vector_index_fields sparse 벡터 인덱스를 담은 컬럼 이름들(콤마 구분). ArrayType(IntegerType)
sparse_vector_value_fields sparse 벡터 값을 담은 컬럼 이름들(콤마 구분). ArrayType(FloatType)
sparse_vector_names 컬렉션 내 sparse 벡터 이름들(콤마 구분). -
multi_vector_fields multi-vector 값을 담은 컬럼 이름들(콤마 구분). ArrayType(ArrayType(FloatType))
multi_vector_names 컬렉션 내 multi-vector 이름들(콤마 구분). -
shard_key_selector upsert 중 사용할 커스텀 shard key 이름들(콤마 구분). -
wait 각 배치 upsert의 완료를 기다릴지 여부. true 또는 false. 기본값: true. -

더 자세한 내용은 Qdrant-Spark GitHub 저장소를 꼭 확인해 보세요. Apache Spark 가이드는 여기에서 볼 수 있어요. 즐거운 데이터 처리 되세요!

더 알아보기 (Learn more)