Data Source V2
Data Source V2 (데이터 소스 V2)
Data Source V2(DSV2)는 외부 데이터 시스템을 Spark에 통합하기 위한 확장 가능한 API예요. org.apache.spark.sql.connector 패키지의 Java 인터페이스 집합으로, 커넥터가 Spark의 쿼리 계획과 실행에 끼어들 수 있게 해줘요. 커넥터는 특정 mix-in 인터페이스를 구현해 기능과 최적화(필터 푸시다운, 컬럼형 읽기, 카탈로그 지원 등)를 선택해서 쓸 수 있어요. 그래서 최소한의 커넥터도 기능을 점진적으로 추가할 수 있죠. 주요 사용자로는 다음이 있어요.
- JDBC 데이터 소스(Spark 내장)
- Apache Iceberg
- Delta Lake
- Lance
이전 Data Source V1 API와 비교해 DSV2가 제공하는 것은 다음과 같아요.
출처: Data Source V2
본문
| 기능 | 설명 |
|---|---|
| Java API | 커넥터 인터페이스가 순수 Java(org.apache.spark.sql.connector)라서 DSV1이 요구하던 Scala 의존성을 없앴어요. Python으로만 작성된 가벼운 커넥터용 래퍼인 Python Data Source API(pyspark.sql.datasource)도 제공돼요. |
| 카탈로그 통합 | 커넥터가 네임스페이스, 테이블, 뷰, 함수를 Spark SQL을 통해 네이티브하게 노출할 수 있어요. |
| 연산자 푸시다운 | 커넥터가 푸시다운된 필터, 필요한 컬럼, 집계, limit, offset 등을 받아들일 수 있어요. |
| 파티셔닝·정렬 보고 | 커넥터가 데이터의 물리적 배치를 보고해서 Spark가 불필요한 shuffle과 sort를 피할 수 있어요. |
| 요청된 분포·정렬 | 쓰기 커넥터가 Spark에게 입력 데이터를 쓰기 전 repartition·정렬하도록 요청해서 클러스터링이나 Z-정렬 같은 최적화된 데이터 레이아웃을 가능하게 해요. |
| 컬럼형 읽기 | 커넥터가 벡터화 처리를 위해 데이터를 컬럼형 배치로 반환할 수 있어요. |
| 행 수준 DML | 커넥터가 전용 인터페이스로 DELETE, UPDATE, MERGE INTO 연산을 네이티브하게 지원할 수 있어요. |
| 스트리밍 지원 | 통합된 Table 추상화가 같은 인터페이스로 배치, 마이크로 배치, 연속 처리를 지원해요. |
진입점(Entry Points)
데이터 소스를 Spark에 연결하는 방법은 두 가지예요.
| 진입점 | 사용 사례 |
|---|---|
| TableProvider | 더 간단한 진입점. 카탈로그를 통하지 않고 옵션(예: 파일 경로나 Kafka 토픽)으로 테이블을 식별하는 소스용. SupportsCatalogOptions를 구현하면 세션 카탈로그를 통한 DDL에도 참여할 수 있어요. |
| CatalogPlugin | 보통 네임스페이스·테이블·선택적으로 뷰·함수의 자체 카탈로그를 관리하는 외부 데이터 소스(Iceberg, Delta Lake 등)가 사용해요. spark.sql.catalog.<name>=com.example.MyCatalog로 등록해요. |
TableProvider
TableProvider는 옵션 집합이 주어지면 Table을 반환해요. 구현체는 public no-arg 생성자가 있어야 해요.
핵심 메서드는 다음과 같아요.
| 메서드 | 설명 |
|---|---|
| inferSchema(options) | 주어진 옵션에서 테이블의 스키마를 추론해요 |
| inferPartitioning(options) | 선택적으로 테이블의 파티셔닝을 추론해요 |
| getTable(schema, partitioning, properties) | 해석된 스키마와 파티셔닝에 대한 Table을 반환해요 |
TableProvider는 또한 SupportsCatalogOptions를 구현해 사용자가 제공한 옵션에서 카탈로그 식별자를 추출해 CREATE TABLE 같은 DDL에 참여할 수 있어요. 이렇게 옵션 기반 소스를 세션 카탈로그에 연결하는데, 이것이 내장 파일 데이터 소스(Parquet, ORC 등)가 테이블 생성을 지원하는 방식이에요.
CatalogPlugin
CatalogPlugin은 카탈로그 구현을 위한 마커 인터페이스예요. 인스턴스화 후 Spark은 카탈로그 이름과 spark.sql.catalog.<name>. 접두어를 공유하는 모든 구성 속성을 사용해 initialize(name, options)를 호출해요.
카탈로그는 추가 카탈로그 인터페이스를 mix-in해 기능을 추가해요.
| 인터페이스 | 기능 |
|---|---|
| TableCatalog | 테이블 나열, 로드, 생성, 변경, 삭제 |
| StagingTableCatalog | 원자적 create-table-as-select / replace-table-as-select |
| SupportsNamespaces | 네임스페이스 생성, 변경, 삭제, 나열 |
| ViewCatalog | 뷰 나열, 로드, 생성, 변경, 삭제(진행 중 — 아직 쿼리 해석에 통합되지 않음) |
| FunctionCatalog | 함수 나열과 로드 |
| ProcedureCatalog | 저장 프로시저 로드와 나열 |
CatalogExtension은 Spark의 내장 세션 카탈로그를 감싸는 특별한 변형으로, 기본 구현에 위임하면서 사용자 지정 동작을 추가하는 데 쓸 수 있어요.
카탈로그 인터페이스
아래 인터페이스들은 CatalogPlugin이 구현해 Spark SQL을 통해 특정 범주의 메타데이터 연산을 노출하기 위한 mix-in이에요.
TableCatalog
TableCatalog은 CatalogPlugin을 확장하며 Table 수명주기 관리를 위한 메서드를 제공해요.
| 메서드 | 설명 |
|---|---|
| listTables(namespace) | 네임스페이스의 테이블을 나열해요 |
| loadTable(ident) | 식별자로 테이블을 로드해요 |
| createTable(ident, columns, partitions, properties) | 새 테이블을 만들어요 |
| alterTable(ident, changes...) | 테이블의 스키마, 속성, 제약을 변경해요 |
| dropTable(ident) | 테이블을 삭제해요 |
StagingTableCatalog
StagingTableCatalog은 TableCatalog를 확장하며 원자적 CREATE TABLE AS SELECT와 REPLACE TABLE AS SELECT를 가능하게 해요. StagedTable을 반환하는데, commitStagedChanges()가 호출될 때만 그 변경이 보이게 돼요. 쓰기가 실패하면 카탈로그는 변하지 않아요.
SupportsNamespaces
SupportsNamespaces는 네임스페이스(데이터베이스/스키마) 관리를 추가해요.
| 메서드 | 설명 |
|---|---|
| listNamespaces() | 하위 네임스페이스를 나열해요 |
| createNamespace(namespace, metadata) | 네임스페이스를 만들어요 |
| alterNamespace(namespace, changes...) | 네임스페이스 속성을 변경해요 |
| dropNamespace(namespace, cascade) | 네임스페이스를 삭제해요 |
ViewCatalog
참고: ViewCatalog는 진행 중인 작업이에요. 인터페이스는 정의됐지만 아직 Spark의 쿼리 해석이나 계획에 통합되지 않았어요.
ViewCatalog은 CatalogPlugin을 확장하며 뷰 수명주기 관리를 위한 메서드를 제공해요.
| 메서드 | 설명 |
|---|---|
| listViews(namespace) | 네임스페이스의 뷰를 나열해요 |
| loadView(ident) | 식별자로 뷰를 로드해요 |
| createView(viewInfo) | 새 뷰를 만들어요 |
| replaceView(viewInfo, orCreate) | 뷰를 교체(또는 생성)해요 |
| alterView(ident, changes...) | 뷰의 속성이나 스키마를 변경해요 |
| dropView(ident) | 뷰를 삭제해요 |
| renameView(oldIdent, newIdent) | 뷰 이름을 바꿔요 |
FunctionCatalog
FunctionCatalog은 사용자 정의 함수 관리를 추가해요.
| 메서드 | 설명 |
|---|---|
| listFunctions(namespace) | 네임스페이스의 함수를 나열해요 |
| loadFunction(ident) | 식별자로 UnboundFunction을 로드해요 |
ProcedureCatalog
ProcedureCatalog은 저장 프로시저 지원을 추가해요.
| 메서드 | 설명 |
|---|---|
| listProcedures(namespace) | 네임스페이스의 프로시저를 나열해요 |
| loadProcedure(ident) | 식별자로 UnboundProcedure를 로드해요 |
프로시저는 CALL catalog.procedure(args)로 호출돼요.
Table
Table은 논리적 데이터셋을 나타내는 중심 추상화예요. 예를 들어 Parquet 파일 디렉터리, Kafka 토픽, 외부 메타스토어가 관리하는 테이블 같은 것이죠. Table은 mix-in 인터페이스를 통해 읽기·쓰기 능력을 얻어요.
Table이 제공하는 것은 다음과 같아요.
| 메서드 | 설명 |
|---|---|
| name() | 테이블의 사람이 읽을 수 있는 식별자 |
| columns() | 테이블의 컬럼(더 이상 사용되지 않는 schema() 메서드를 대체) |
| partitioning() | Transform 배열로 표현된 물리적 파티셔닝 |
| properties() | 테이블 속성의 문자열 맵 |
| capabilities() | 테이블이 지원하는 것을 선언하는 TableCapability 값 집합 |
읽기·쓰기 Mix-in
Table은 mix-in 인터페이스를 구현해 읽기·쓰기 능력을 얻어요.
- **
SupportsRead**는newScanBuilder(options)를 추가하며, 배치 또는 스트리밍 읽기용ScanBuilder를 반환해요. Read Path 참고. - **
SupportsWrite**는newWriteBuilder(info)를 추가하며, 배치 또는 스트리밍 쓰기용WriteBuilder를 반환해요. Write Path 참고. - **
SupportsRowLevelOperations**는 읽고 다시 쓰는(read-and-rewrite) 주기를 통해DELETE,UPDATE,MERGE INTO를 가능하게 해요. Row-Level DML 참고.
추가 mix-in은 더 많은 기능을 가능하게 해요.
| Mix-in | 기능 |
|---|---|
| SupportsDelete / SupportsDeleteV2 | 필터 기반 행 삭제 |
| TruncatableTable | TRUNCATE TABLE |
| SupportsPartitionManagement | 파티션 DDL(ADD/DROP/RENAME PARTITION) |
| SupportsMetadataColumns | 숨은 메타데이터 컬럼(예: 파일 이름, 행 위치) 노출 |
읽기 경로(Read Path)
읽기 경로는 논리적 계획과 물리적 실행을 분리하는 빌더 패턴을 따르는 구조예요.
SupportsRead.newScanBuilder(options)
└─▸ ScanBuilder (logical: pushdown negotiation)
└─▸ Scan (logical: read schema, description)
└─▸ Batch (physical: partitions + reader factory)
├─ InputPartition[]
└─ PartitionReaderFactory
└─▸ PartitionReader (per-task I/O)
ScanBuilder
ScanBuilder는 읽기를 구성하는 시작점이에요. Spark은 build()를 호출해 최종 Scan을 얻어요.
build()를 호출하기 전에 Spark은 mix-in 인터페이스를 확인해 스캔 빌더와 연산자 푸시다운을 협상해요. 푸시다운 순서는 다음과 같아요.
- Sample(
SupportsPushDownSample) - Filter(
SupportsPushDownFilters/SupportsPushDownV2Filters) - Aggregate(
SupportsPushDownAggregates) - Limit / Top-N(
SupportsPushDownLimit/SupportsPushDownTopN) - Offset(
SupportsPushDownOffset) - Column pruning(
SupportsPushDownRequiredColumns)
각 mix-in 인터페이스는 Spark이 관련 연산자와 함께 호출하는 메서드를 가져요. 구현은 처리할 수 있는 연산자를 반환하고, Spark이 나머지 연산자를 직접 적용해요.
Scan
Scan은 구성된 데이터 소스 읽기의 논리적 표현이에요. 제공하는 것은 다음과 같아요.
| 메서드 | 설명 |
|---|---|
| readSchema() | 컬럼 프루닝이나 푸시다운 후의 실제 스키마 |
| toBatch() | 배치 실행용 Batch 반환 |
| toMicroBatchStream(checkpointLocation) | 스트리밍용 MicroBatchStream 반환 |
| toContinuousStream(checkpointLocation) | 연속 처리를 위한 ContinuousStream 반환 |
구현체는 자신의 Table이 선언한 TableCapability에 해당하는 메서드를 반드시 오버라이드해야 해요.
Batch
Batch는 배치 읽기의 물리적 표현이에요. 두 가지 메서드가 있어요.
| 메서드 | 설명 |
|---|---|
| planInputPartitions() | InputPartition 객체의 배열을 반환해요. 각 파티션은 Spark 태스크 하나에 매핑돼요 |
| createReaderFactory() | 직렬화되어 실행자로 보내지는 PartitionReaderFactory를 반환해요 |
InputPartition과 PartitionReader
InputPartition은 데이터 분할을 나타내는 직렬화 가능한 핸들이에요. 데이터 지역성을 위해 선택적으로 preferredLocations()를 선언할 수 있어요.
PartitionReaderFactory는 실행자로 직렬화되어 각 입력 파티션에 대해 PartitionReader를 만들어요. 리더는 행(또는 컬럼형 소스의 경우 ColumnarBatch 인스턴스)을 반복하고 소비 후 닫혀요.
Scan Mix-in
Scan 구현은 mix-in 인터페이스를 구현해 추가 최적화를 선택할 수 있어요.
| Mix-in | 기능 |
|---|---|
| SupportsReportPartitioning | 출력이 어떻게 파티셔닝되는지 보고해 Spark이 불필요한 shuffle을 피하게 해요(예: Storage Partition Join) |
| SupportsReportOrdering | 출력의 정렬 순서를 보고해 Spark이 중복 정렬을 건너뛰게 해요 |
| SupportsReportStatistics | 비용 기반 옵티마이저를 위한 테이블·컬럼 수준 통계를 보고해요 |
| SupportsRuntimeV2Filtering | 실행 시간에 추가 필터 값을 받아 동적 파티션 프루닝과 행 수준 DML 재작성 범위 축소에 사용해요 |
쓰기 경로(Write Path)
쓰기 경로는 읽기 경로처럼 논리적 구성과 물리적 실행을 분리하는 빌더 패턴을 따르는 구조예요.
SupportsWrite.newWriteBuilder(info)
└─▸ WriteBuilder (logical: mode selection)
└─▸ Write (logical: description, metrics)
└─▸ BatchWrite (physical: commit protocol)
├─ DataWriterFactory
│ └─▸ DataWriter (per-task I/O)
├─ commit(messages[])
└─ abort(messages[])
WriteBuilder
WriteBuilder은 쓰기를 구성하는 시작점이에요. build()를 호출하면 논리적 Write 객체를 반환해요.
쓰기 모드는 빌더에 추가 인터페이스를 mix-in해 구성해요.
| Mix-in | 모드 |
|---|---|
| SupportsTruncate | 쓰기 전에 테이블을 비워요 |
| SupportsOverwriteV2 | 필터 표현식과 일치하는 데이터를 덮어써요 |
| SupportsDynamicOverwrite | 파티션을 동적으로 덮어써요 |
Write
Write는 구성된 쓰기의 논리적 표현이에요. Scan과 비슷하게 물리적 계층으로 연결돼요.
| 메서드 | 설명 |
|---|---|
| toBatch() | 배치 실행용 BatchWrite 반환 |
| toStreaming() | 스트리밍 실행용 StreamingWrite 반환 |
BatchWrite
BatchWrite는 배치 쓰기를 위한 2단계 커밋(two-phase commit) 프로토콜을 정의해요.
createBatchWriterFactory(info)— 직렬화되어 실행자로 보내지는DataWriterFactory를 만들어요.- 각 실행자에서 팩토리는 파티션마다
DataWriter를 만들어요. 모든 행이 성공적으로 쓰이면commit()을, 그렇지 않으면abort()를 호출해요. - 모든 태스크가 끝나면 드라이버는 (모든 태스크가 성공했다면)
commit(messages[])또는 (어느 태스크라도 실패했다면)abort(messages[])를 호출해요.
개별 태스크가 쓴 데이터는 드라이버 수준의 commit이 성공하기 전까지는 리더에게 보이면 안 돼요.
분포·정렬 요구사항
Write 구현은 또한 RequiresDistributionAndOrdering를 구현해 Spark에게 쓰기 전에 입력 데이터가 어떻게 분포·정렬돼야 하는지 알려줄 수 있어요. Spark은 이 요구사항을 충족하도록 필요에 따라 shuffle과 sort 노드를 삽입해요.
행 수준 DML
DSV2는 커넥터가 DELETE, UPDATE, MERGE INTO 문을 지원할 수 있게 하는 인터페이스를 제공해요.
필터 기반 삭제
DML의 가장 간단한 형태는 필터 기반 삭제예요. 술어와 일치하는 전체 행 그룹이 개별 레코드를 다시 쓰지 않고 제거돼요.
SupportsDeleteV2(및 오래된 SupportsDelete)는 이 용도의 Table mix-in이에요. 핵심 메서드는 다음과 같아요.
| 메서드 | 설명 |
|---|---|
| canDeleteWhere(predicates) | 주어진 술어와 일치하는 행을 데이터 소스가 효율적으로 삭제할 수 있는지 반환해요. false를 반환하면 Spark은 행 수준 재작성으로 폴백해요(아래 참고) |
| deleteWhere(predicates) | 술어와 일치하는 모든 행을 삭제해요 |
이 접근법은 데이터를 다시 쓰지 않고 전체 파티션이나 파일을 버릴 수 있는 데이터 소스에 잘 맞아요.
행 수준 연산
UPDATE, MERGE INTO, 또는 단순 필터로 처리할 수 없는 DELETE 같은 더 복잡한 DML에는 커넥터가 SupportsRowLevelOperations를 구현해요. 이 인터페이스는 읽고 다시 쓰는 주기를 조정하는 RowLevelOperation을 반환해요.
SupportsRowLevelOperations.newRowLevelOperationBuilder(info)
└─▸ RowLevelOperation
├─ newScanBuilder(options) → reads affected rows
└─ newWriteBuilder(info) → writes rewritten data
데이터 소스는 두 범주로 나뉘어요.
- 델타 기반(Delta-based) 소스(때로 merge-on-read라고 함)는 개별 행 변경(삽입, 갱신, 삭제)의 스트림을 처리할 수 있어요. 스캔은 변경 중인 행만 만들면 돼요.
- 그룹 기반(Group-based) 소스(때로 copy-on-write라고 함)는 전체 행 그룹(예: 파일이나 파티션)을 교체해요. 스캔은 변경되지 않은 행을 포함해 각 영향받은 그룹의 모든 행을 반환해야 해서, 데이터 소스가 수정 사항을 적용해 그룹을 다시 쓸 수 있어요.
표현식(Expressions)
org.apache.spark.sql.connector.expressions 패키지는 DSV2 API 전반에서 쓰는 중립적인 표현식 표현을 제공해요.
- Transforms(
IdentityTransform,BucketTransform,YearsTransform등)는 테이블 파티셔닝 방식을 나타내요. - 정렬 순서(
SortOrder)는 정렬 요구사항을 나타내요. - 필터 술어(
Predicate와 하위 클래스)는 푸시다운된 필터 조건을 나타내요. - 집계(Aggregates) 는 푸시다운된 집계 함수를 나타내요.
이 표현식 타입들은 Spark의 내부 Catalyst 표현식과 독립적이어서 커넥터가 Spark 내부에 대한 의존성을 피할 수 있어요.
스트리밍
동일한 Table과 Scan 추상화가 스트리밍 쿼리를 지원해요. MICRO_BATCH_READ나 CONTINUOUS_READ 능력을 선언한 Table은 다음을 통해 스트리밍 읽기를 제공해요.
| 메서드 | 설명 |
|---|---|
| Scan.toMicroBatchStream(checkpointLocation) | 마이크로 배치로 데이터를 읽고 offset으로 진행 상황을 추적하는 MicroBatchStream을 반환해요 |
| Scan.toContinuousStream(checkpointLocation) | 저지연 연속 처리를 위한 ContinuousStream을 반환해요 |
STREAMING_WRITE를 선언한 Table은 Write.toStreaming()을 통해 스트리밍 쓰기를 지원하며, 이것은 StreamingWrite를 반환해요.
더 읽을거리(Further Reading)
- 전체 인터페이스 참조는 API 문서(Javadoc).
- 내장 데이터 소스(DSV1)에 대한 사용자 가이드는 Data Sources.
- Python으로만 작성된 가벼운 커넥터는 Python Data Source API.
- DSV2 파티셔닝 보고가 조인 최적화를 가능하게 하는 방법은 Storage Partition Join.
더 알아보기 (Learn more)
- Data Sources — 내장 데이터 소스 사용자 가이드.
- JDBC 데이터 소스 — DSV2 기반 내장 JDBC 소스.
- Spark SQL 프로그래밍 가이드 — DataFrame과 SQL 전반.