Apache Spark 연동
Apache Spark 연동 (spark)
Spark는 빅데이터 처리와 분석을 위해 설계된 분산 컴퓨팅 프레임워크예요. Qdrant-Spark 커넥터를 사용하면 Spark에서 Qdrant를 storage destination(저장 목적지)으로 사용할 수 있게 됩니다.
설치 (Installation)
커넥터를 Spark 환경에 통합하려면, 아래 나열된 소스 중 하나에서 JAR 파일을 받으면 돼요.
GitHub Releases
모든 필수 의존성이 포함된 패키징된 jar 파일은 여기에서 찾을 수 있어요.
소스에서 빌드하기
jar을 소스에서 빌드하려면 JDK@8과 Maven이 설치되어 있어야 해요. 요구사항이 충족되면 프로젝트 루트에서 다음 명령을 실행하세요.
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_field와vector_name옵션은 하위 호환성을 위해 유지되고 있어요. named 벡터의 경우 아래와 같이vector_fields와vector_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 가이드는 여기에서 볼 수 있어요. 즐거운 데이터 처리 되세요!