스파크 쓰기

스파크 쓰기 (Spark Writes)

이 문서에서는 스파크에서 아이스버그 테이블에 데이터를 쓰는 방법을 알려드릴게요. SQL의 INSERT INTO, MERGE INTO, INSERT OVERWRITE부터 DataFrameWriterV2 API, 브랜치 쓰기, 분포 모드, 파일 크기 제어까지 폭넓게 살펴볼게요. 아이스버그는 스파크의 DataSourceV2 API를 사용하며, 행 수준 업데이트 같은 고급 기능을 지원해요.

출처: 문서

본문

스파크에서 아이스버그를 사용하려면 먼저 스파크 카탈로그를 구성해요.

일부 계획(plan)은 아이스버그 SQL 확장을 사용할 때만 사용할 수 있어요.

아이스버그는 데이터 소스와 카탈로그 구현에 아파치 스파크의 DataSourceV2 API를 사용해요. Spark DSv2는 스파크 버전마다 지원 수준이 다른 진화하는 API예요.

기능 지원 Spark 참고
SQL insert into ✔️ ⚠ spark.sql.storeAssignmentPolicy=ANSI 필요 (Spark 3.0부터 기본값)
SQL merge into ✔️ ⚠ 아이스버그 스파크 확장 필요
SQL insert overwrite ✔️ ⚠ spark.sql.storeAssignmentPolicy=ANSI 필요 (Spark 3.0부터 기본값)
SQL delete from ✔️ ⚠ 행 수준 삭제는 아이스버그 스파크 확장 필요
SQL update ✔️ ⚠ 아이스버그 스파크 확장 필요
DataFrame append ✔️
DataFrame overwrite ✔️
DataFrame CTAS and RTAS ✔️ ⚠ DSv2 API 필요
DataFrame merge into ✔️ ⚠ DSv2 API 필요 (Spark 4.0 이상)

SQL로 쓰기 (Writing with SQL)

스파크는 SQL INSERT INTO, MERGE INTO, INSERT OVERWRITE와 새 DataFrameWriterV2 API를 지원해요.

INSERT INTO

테이블에 새 데이터를 추가하려면 INSERT INTO를 사용해요.

INSERT INTO prod.db.table VALUES (1, 'a'), (2, 'b')

INSERT INTO prod.db.table SELECT ...

MERGE INTO

스파크는 행 수준 업데이트를 표현할 수 있는 MERGE INTO 쿼리를 지원해요.

아이스버그는 업데이트해야 하는 행을 포함한 데이터 파일을 overwrite 커밋으로 재작성해서 MERGE INTO를 지원해요.

MERGE INTO는 INSERT OVERWRITE 대신 권장돼요. 아이스버그가 영향받은 데이터 파일만 교체할 수 있고, 동적 overwrite가 덮어쓰는 데이터는 테이블 파티셔닝이 바뀌면 달라질 수 있기 때문이에요.

MERGE INTO 구문 (MERGE INTO syntax)

MERGE INTO는 target 테이블이라고 부르는 테이블을, source라고 부르는 다른 쿼리의 업데이트 집합을 사용해서 업데이트해요. target 테이블의 행에 대한 업데이트는 조인 조건 같은 ON 절을 사용해서 찾아요.

MERGE INTO prod.db.target t   -- a target table
USING (SELECT ...) s          -- the source updates
ON t.id = s.id                -- condition to find updates for target rows
WHEN ...                      -- updates

target 테이블의 행에 대한 업데이트는 WHEN MATCHED ... THEN ...으로 나열돼요. 각 매치가 언제 적용돼야 하는지를 결정하는 조건으로 여러 MATCHED 절을 추가할 수 있어요. 첫 번째로 일치하는 표현식이 사용돼요.

WHEN MATCHED AND s.op = 'delete' THEN DELETE
WHEN MATCHED AND t.count IS NULL AND s.op = 'increment' THEN UPDATE SET t.count = 0
WHEN MATCHED AND s.op = 'increment' THEN UPDATE SET t.count = t.count + 1

일치하지 않는 소스 행(업데이트)은 삽입할 수 있어요.

WHEN NOT MATCHED THEN INSERT *

삽입은 추가 조건도 지원해요.

WHEN NOT MATCHED AND s.event_time > still_valid_threshold THEN INSERT (id, count) VALUES (s.id, 1)

소스 데이터의 레코드 하나만 target 테이블의 주어진 행을 업데이트할 수 있고, 그렇지 않으면 오류가 발생해요.

Spark 3.5는 소스 데이터에 없는 행을 업데이트하거나 삭제하는 WHEN NOT MATCHED BY SOURCE ... THEN ... 지원을 추가했어요.

WHEN NOT MATCHED BY SOURCE THEN UPDATE SET status = 'invalid'

스냅샷 요약 (Snapshot summary)

MERGE INTO 커밋 후 스냅샷 요약에는 다음 필드들이 포함될 수 있어요. 각 값은 음수가 아닌 개수의 문자열 형태예요. 값이 알 수 없을 때(예: 스파크가 보고하지 않음)는 필드가 생략돼요.

정보 (Info)

Spark 4.1 이상에서만 사용할 수 있어요.

필드 설명
spark.merge-into.num-target-rows-copied 어떤 작업과도 일치하지 않아 수정되지 않고 복사된 target 행 수
spark.merge-into.num-target-rows-deleted 삭제된 target 행 수
spark.merge-into.num-target-rows-updated 업데이트된 target 행 수
spark.merge-into.num-target-rows-inserted 삽입된 target 행 수
spark.merge-into.num-target-rows-matched-updated MATCHED 절로 업데이트된 target 행 수
spark.merge-into.num-target-rows-matched-deleted MATCHED 절로 삭제된 target 행 수
spark.merge-into.num-target-rows-not-matched-by-source-updated NOT MATCHED BY SOURCE 절로 업데이트된 target 행 수
spark.merge-into.num-target-rows-not-matched-by-source-deleted NOT MATCHED BY SOURCE 절로 삭제된 target 행 수

INSERT OVERWRITE

INSERT OVERWRITE는 테이블의 데이터를 쿼리 결과로 교체할 수 있어요. Overwrite는 아이스버그 테이블에서 원자적 연산이에요.

INSERT OVERWRITE가 교체할 파티션은 스파크의 파티션 overwrite 모드와 테이블의 파티셔닝에 따라 달라져요. MERGE INTO는 영향받은 데이터 파일만 재작성하고 동작을 이해하기 더 쉬우므로, INSERT OVERWRITE 대신 권장돼요.

Overwrite 동작 (Overwrite behavior)

스파크의 기본 overwrite 모드는 static이지만, 아이스버그 테이블에 쓸 때는 dynamic overwrite 모드가 권장돼요. Static overwrite 모드는 PARTITION 절을 필터로 변환해서 테이블에서 어떤 파티션을 덮어쓸지 결정하지만, PARTITION 절은 테이블 컬럼만 참조할 수 있어요.

Dynamic overwrite 모드는 spark.sql.sources.partitionOverwriteMode=dynamic을 설정해서 구성돼요.

dynamic과 static overwrite의 동작을 보여주기 위해, 다음 DDL로 정의된 logs 테이블을 고려해볼게요.

CREATE TABLE prod.my_app.logs (
    uuid string NOT NULL,
    level string NOT NULL,
    ts timestamp NOT NULL,
    message string)
USING iceberg
PARTITIONED BY (level, hours(ts))

Dynamic overwrite

스파크의 overwrite 모드가 dynamic이면, SELECT 쿼리가 만든 행을 가진 파티션이 교체돼요.

예를 들어 이 쿼리는 예시 logs 테이블에서 중복 로그 이벤트를 제거해요.

INSERT OVERWRITE prod.my_app.logs
SELECT uuid, first(level), first(ts), first(message)
FROM prod.my_app.logs
WHERE cast(ts as date) = '2020-07-01'
GROUP BY uuid

dynamic 모드에서 이것은 SELECT 결과에 행이 있는 모든 파티션을 교체해요. 모든 행의 날짜가 7월 1일로 제한되므로 그 날짜의 시간(hour)만 교체돼요.

Static overwrite

스파크의 overwrite 모드가 static이면, PARTITION 절이 테이블에서 삭제하는 데 사용되는 필터로 변환돼요. PARTITION 절이 생략되면 모든 파티션이 교체돼요.

위 쿼리에는 PARTITION 절이 없으므로 static 모드로 실행하면 테이블의 기존 모든 행을 삭제하지만, 7월 1일의 로그만 쓰게 돼요.

로드된 파티션만 덮어쓰려면 SELECT 쿼리 필터와 정렬되는 PARTITION 절을 추가해요.

INSERT OVERWRITE prod.my_app.logs
PARTITION (level = 'INFO')
SELECT uuid, first(level), first(ts), first(message)
FROM prod.my_app.logs
WHERE level = 'INFO'
GROUP BY uuid

이 모드는 PARTITION 절이 테이블 컬럼만 참조할 수 있고 숨겨진 파티션은 참조할 수 없기 때문에, dynamic 예시 쿼리처럼 시간 단위 파티션을 교체할 수 없다는 점에 유의해주세요.

DELETE FROM

스파크는 테이블에서 데이터를 제거하는 DELETE FROM 쿼리를 지원해요.

Delete 쿼리는 삭제할 행을 일치시키는 필터를 받아요.

DELETE FROM prod.db.table
WHERE ts >= '2020-05-01 00:00:00' and ts < '2020-06-01 00:00:00'

DELETE FROM prod.db.all_events
WHERE session_time < (SELECT min(session_time) FROM prod.db.good_events)

DELETE FROM prod.db.orders AS t1
WHERE EXISTS (SELECT oid FROM prod.db.returned_orders WHERE t1.oid = oid)

삭제 필터가 테이블의 전체 파티션과 일치하면 아이스버그는 메타데이터 전용 삭제(metadata-only delete)를 수행해요. 필터가 테이블의 개별 행과 일치하면 아이스버그는 영향받은 데이터 파일만 재작성해요.

UPDATE

Update 쿼리는 업데이트할 행을 일치시키는 필터를 받아요.

UPDATE prod.db.table
SET c1 = 'update_c1', c2 = 'update_c2'
WHERE ts >= '2020-05-01 00:00:00' and ts < '2020-06-01 00:00:00'

UPDATE prod.db.all_events
SET session_time = 0, ignored = true
WHERE session_time < (SELECT min(session_time) FROM prod.db.good_events)

UPDATE prod.db.orders AS t1
SET order_status = 'returned'
WHERE EXISTS (SELECT oid FROM prod.db.returned_orders WHERE t1.oid = oid)

들어오는 데이터에 기반한 더 복잡한 행 수준 업데이트는 MERGE INTO 섹션을 참고해주세요.

브랜치 쓰기 (Writing to Branches)

쓰기를 수행하기 전에 브랜치가 존재해야 해요. 연산은 브랜치가 없으면 만들지 않아요. 브랜치는 스파크 DDL로 만들 수 있어요.

정보 (Info)

참고: 브랜치에 쓸 때 검증에는 테이블의 현재 스키마가 사용돼요.

SQL을 통한 브랜치 쓰기 (Via SQL)

브랜치 쓰기는 연산에서 branch_yourBranch 브랜치 식별자를 제공해서 수행할 수 있어요.

브랜치 쓰기는 spark.wap.branch 구성을 지정해서 write-audit-publish(WAP) 워크플로우의 일부로도 수행할 수 있어요. WAP 브랜치와 브랜치 식별자는 둘 다 지정할 수 없어요.

-- INSERT (1,' a') (2, 'b') into the audit branch.
INSERT INTO prod.db.table.branch_audit VALUES (1, 'a'), (2, 'b');

-- MERGE INTO audit branch
MERGE INTO prod.db.table.branch_audit t 
USING (SELECT ...) s        
ON t.id = s.id          
WHEN ...

-- UPDATE audit branch
UPDATE prod.db.table.branch_audit AS t1
SET val = 'c'

-- DELETE FROM audit branch
DELETE FROM prod.db.table.branch_audit WHERE id = 2;

-- WAP Branch write
SET spark.wap.branch = audit-branch
INSERT INTO prod.db.table VALUES (3, 'c');

DataFrames를 통한 브랜치 쓰기 (Via DataFrames)

DataFrames를 통한 브랜치 쓰기는 연산에서 branch_yourBranch 브랜치 식별자를 제공해서 수행할 수 있어요.

// To insert into `audit` branch
val data: DataFrame = ...
data.writeTo("prod.db.table.branch_audit").append()
// To overwrite `audit` branch
val data: DataFrame = ...
data.writeTo("prod.db.table.branch_audit").overwritePartitions()

DataFrames로 쓰기 (Writing with DataFrames)

스파크는 데이터 프레임을 사용해 테이블에 쓰는 새 DataFrameWriterV2 API를 도입했어요. v2 API는 여러 이유로 권장돼요.

  • CTAS, RTAS, 필터로 overwrite가 지원돼요
  • 모든 연산이 이름으로 테이블에 컬럼을 일관되게 써요
  • partitionedBy에서 숨겨진 파티션 표현식이 지원돼요
  • overwrite 동작이 명시적이에요. dynamic이거나 사용자 제공 필터에 의한 방식이에요
  • 각 연산의 동작이 SQL 문과 대응돼요. df.writeTo(t).create()는 CREATE TABLE AS SELECT와 동등하고, df.writeTo(t).replace()는 REPLACE TABLE AS SELECT와 동등하며, df.writeTo(t).append()는 INSERT INTO와 동등하고, df.writeTo(t).overwritePartitions()는 dynamic INSERT OVERWRITE와 동등해요

v1 DataFrame 쓰기 API도 여전히 지원되지만 권장되지는 않아요.

위험 (Danger)

스파크에서 v1 DataFrame API로 쓸 때는 카탈로그로 테이블을 로드하려면 saveAsTable이나 insertInto를 사용해요. format("iceberg")를 사용하면 쿼리가 사용하는 테이블을 자동으로 갱신하지 않는 격리된 테이블 참조를 로드해요.

데이터 추가 (Appending data)

데이터프레임을 아이스버그 테이블에 추가하려면 append를 사용해요.

val data: DataFrame = ...
data.writeTo("prod.db.table").append()

데이터 덮어쓰기 (Overwriting data)

파티션을 동적으로 덮어쓰려면 overwritePartitions()를 사용해요.

val data: DataFrame = ...
data.writeTo("prod.db.table").overwritePartitions()

파티션을 명시적으로 덮어쓰려면 overwrite를 사용해서 필터를 제공해요.

data.writeTo("prod.db.table").overwrite($"level" === "INFO")

테이블 만들기 (Creating tables)

CTAS나 RTAS를 실행하려면 create, replace, createOrReplace 연산을 사용해요.

val data: DataFrame = ...
data.writeTo("prod.db.table").create()

기본 스파크 카탈로그(spark_catalog)를 아이스버그의 SparkSessionCatalog로 교체했다면 다음을 수행해요.

val data: DataFrame = ...
data.writeTo("db.table").using("iceberg").create()

create와 replace 연산은 partitionedBy, tableProperty 같은 테이블 구성 메서드를 지원해요.

data.writeTo("prod.db.table")
    .tableProperty("write.format.default", "orc")
    .partitionedBy($"level", days($"ts"))
    .createOrReplace()

아이스버그 테이블 위치는 location 테이블 속성으로도 지정할 수 있어요.

data.writeTo("prod.db.table")
    .tableProperty("location", "/path/to/location")
    .createOrReplace()

데이터 병합 (Merging data)

Spark 4.0은 DataFrameWriterV2 API로 MERGE INTO 쿼리를 수행하는 지원을 추가했어요.

MERGE INTO 쿼리는 소스(이 경우 DataFrame)의 업데이트 집합을 사용해서 target 테이블을 업데이트해요.

val source: DataFrame = ...                               // e.g., read from a table, "source"
source.mergeInto("target", $"source.id" === $"target.id") // second argument is the ON condition
    .whenMatched($"target.id" === 1)                      // argument is the additional condition
    .updateAll()                                          // UPDATE SET *
    .whenMatched($"target.id" === 2)
    .delete()
    .whenNotMatched()
    .insertAll()                                          // INSERT *
    .whenNotMatchedBySource($"target.id" === 3)
    .update(Map("status" -> lit("invalid")))              // set column name(s) to expression(s)
    .merge()

스키마 병합 (Schema Merge)

삽입하거나 업데이트할 때 아이스버그는 런타임에 스키마 불일치를 해결할 수 있어요. 구성되면 아이스버그는 다음과 같이 자동 스키마 진화를 수행해요.

  • 소스에는 있지만 target 테이블에는 없는 새 컬럼이 있을 때. 새 컬럼이 target 테이블에 추가돼요. 테이블에 이미 있는 모든 행의 컬럼 값은 NULL로 설정돼요.
  • target에는 있지만 소스에는 없는 컬럼이 있을 때. target 컬럼 값은 삽입할 때 NULL로 설정되거나 행을 업데이트할 때 그대로 둬요.

target 테이블은 write.spark.accept-any-schema 속성을 true로 설정해서 어떤 스키마 변경도 받아들이도록 구성해야 해요.

ALTER TABLE prod.db.sample SET TBLPROPERTIES (
  'write.spark.accept-any-schema'='true'
)

작성자는 mergeSchema 옵션을 활성화해야 해요.

data.writeTo("prod.db.sample").option("mergeSchema","true").append()

쓰기 분포 모드 (Writing Distribution Modes)

아이스버그의 기본 스파크 작성자는 각 스파크 작업의 데이터가 파티션 값으로 클러스터링돼야 해요. 이 분포는 쓰기 중에 열려 있는 파일 핸들의 수를 최소화하는 데 필요해요. 기본적으로 Iceberg 1.2.0부터 아이스버그는 스파크가 이 분포에 맞게 쓸 데이터를 미리 정렬하도록 요청해요. 스파크에 대한 요청은 write.distribution-mode 테이블 속성으로 값 hash를 사용해서 수행돼요. 스파크는 3.5.0 이전에는 CTAS/RTAS에서 분포 모드를 존중하지 않아요.

다음 샘플 테이블에 데이터를 쓰는 과정을 살펴볼게요.

CREATE TABLE prod.db.sample (
    id bigint,
    data string,
    category string,
    ts timestamp)
USING iceberg
PARTITIONED BY (days(ts), category)

샘플 테이블에 데이터를 쓰려면 데이터를 days(ts), category로 정렬해야 하는데, 이는 기본 hash 분포가 자동으로 처리해요. 이전에는 수동 정렬이 필요했지만 이제는 그렇지 않아요.

INSERT INTO prod.db.sample
SELECT id, data, category, ts FROM another_table

write.distribution-mode에는 세 가지 옵션이 있어요.

  • none - 이것은 아이스버그의 이전 기본값이에요. 이 모드는 스파크가 자동으로 셔플이나 정렬을 수행하도록 요청하지 않아요. 스파크가 자동으로 아무 작업도 하지 않으므로, 데이터를 파티션 값으로 수동으로 정렬해야 해요. 데이터는 각 스파크 작업 안에서 정렬하거나, 전체 데이터셋 안에서 전역적으로 정렬해야 해요. 전역 정렬은 출력 파일 수를 최소화해요. 스파크 쓰기 fanout 속성을 사용하면 정렬을 피할 수 있지만, 이렇게 하면 각 쓰기 작업이 완료될 때까지 모든 파일 핸들이 열린 채 유지돼요.
  • hash - 이 모드는 새 기본값이고, 스파크가 쓰기 전에 들어오는 쓰기 데이터를 셔플하기 위해 hash 기반 교환을 사용하도록 요청해요. 실제로 이는 각 행이 행의 파티션 값을 기반으로 해시된 다음 그 값에 해당하는 스파크 작업에 배치된다는 뜻이에요. 스파크의 Adaptive Query 계획 때문에 작업의 추가 분할과 병합이 일어날 수 있어요.
  • range - 이 모드는 스파크가 쓰기 전에 데이터를 셔플하기 위해 범위 기반 교환(range based exchange)을 수행하도록 요청해요. 이것은 hash 모드보다 더 비싼 2단계 절차예요. 첫 번째 단계는 파티션과 정렬 컬럼을 기반으로 쓸 데이터를 샘플링해요. 두 번째 단계는 범위 정보를 사용해 입력 데이터를 스파크 작업으로 셔플해요. 각 작업은 입력 데이터의 배타적 범위를 얻는데, 이는 파티션으로 데이터를 클러스터링하고 전역적으로도 정렬해요. hash 분포보다 비싸지만, 정렬된 컬럼이 쿼리 중에 사용되면 전역 정렬은 읽기 성능에 유용할 수 있어요. 이 모드는 테이블이 정렬 순서(sort-order)로 만들어지면 기본으로 사용돼요. 스파크의 Adaptive Query 계획 때문에 작업의 추가 분할과 병합이 일어날 수 있어요.

파일 크기 제어 (Controlling File Sizes)

스파크로 아이스버그에 데이터를 쓸 때, 스파크는 스파크 작업보다 큰 파일을 쓸 수 없고 파일은 아이스버그 파티션 경계를 넘을 수 없다는 점을 기억하는 것이 중요해요. 이는 아이스버그가 항상 write.target-file-size-bytes만큼 커지면 파일을 롤오버하지만, 스파크 작업이 충분히 크지 않으면 그렇게 되지 않는다는 뜻이에요. 디스크에 생성되는 파일의 크기는 압축되지 않은 스파크 행 표현과 달리 디스크의 데이터가 압축되고 컬럼 포맷이므로 스파크 작업보다 훨씬 작아요. 즉 100MB 스파크 작업은 단일 아이스버그 파티션에 쓴다 해도 100MB보다 훨씬 작은 파일을 만들게 돼요. 작업이 여러 파티션에 쓰면 파일은 그보다도 더 작아져요.

각 스파크 작업에 어떤 데이터가 들어갈지 제어하려면 쓰기 분포 모드를 사용하거나 데이터를 수동으로 repartition 해요.

스파크 작업 크기를 조정하려면 스파크의 다양한 Adaptive Query Execution(AQE) 파라미터를 잘 알아야 해요. write.distribution-mode가 none이 아니면, AQE는 교환 중에 스파크 작업의 병합과 분할을 제어해서 spark.sql.adaptive.advisoryPartitionSizeInBytes 크기인 작업을 만들려고 해요. 이 설정은 사용자가 수행한 re-partition이나 정렬에도 영향을 줘요. 이것은 디스크의 컬럼 압축 크기가 아니라 인메모리 스파크 행 크기이므로, 목표 파일 크기보다 큰 값을 지정해야 한다는 점을 다시 기억하는 것이 중요해요. 인메모리 크기와 디스크 크기의 비율은 데이터에 따라 달라져요. 스파크의 향후 작업에서 아이스버그가 쓰기 시점에 이 파라미터를 자동으로 조정해서 write.target-file-size-bytes와 일치시킬 수 있어야 해요.

더 알아보기 (Learn more)