스파크 구조적 스트리밍
스파크 구조적 스트리밍 (Spark Structured Streaming)
아이스버그는 스파크 구조적 스트리밍(Structured Streaming)에서 증분 데이터를 읽고 쓸 수 있어요. 이 문서에서는 스트리밍 읽기와 쓰기 설정, 입력 속도 제한, 비동기 마이크로 배치 계획, 그리고 스트리밍 테이블의 유지보수 방법을 알려드릴게요. 아이스버그는 스파크 엔진과 자주 커밋되는 스트리밍 워크로드에 맞춰 세밀하게 튜닝할 수 있는 여러 옵션을 제공해요.
출처: 문서
본문
아이스버그는 데이터 소스와 카탈로그 구현에 아파치 스파크의 DataSourceV2 API를 사용해요. Spark DSv2는 스파크 버전마다 지원 수준이 다른 진화하는 API예요.
스트리밍 읽기 (Streaming Reads)
아이스버그는 과거 타임스탬프에서 시작하는 스파크 구조적 스트리밍 작업에서 증분 데이터를 처리하는 것을 지원해요.
val df = spark.readStream
.format("iceberg")
.option("stream-from-timestamp", Long.toString(streamStartTimestamp))
.load("database.table_name")
경고 (Warning)
아이스버그는 추가(append) 스냅샷의 데이터만 읽는 것을 지원해요. Overwrite 스냅샷은 처리할 수 없고 기본적으로 예외를 발생시켜요. streaming-skip-overwrite-snapshots=true로 설정하면 overwrite를 무시할 수 있어요. 마찬가지로 delete 스냅샷도 기본적으로 예외를 일으키고, streaming-skip-delete-snapshots=true로 설정하면 delete를 무시할 수 있어요.
입력 속도 제한 (Limit input rate)
DataFrame API에서 마이크로 배치의 크기를 제어하기 위해 아이스버그는 두 가지 읽기 옵션을 지원해요.
- streaming-max-files-per-micro-batch: 매 마이크로 배치마다 처리할 파일의 최대 개수.
- streaming-max-rows-per-micro-batch: 매 마이크로 배치마다 처리할 행 수의 "소프트 최대값". 배치는 항상 다음 미처리 데이터 파일의 모든 행을 포함하지만, 소프트 최대값을 초과하게 된다면 추가 파일은 포함되지 않아요.
두 옵션을 모두 설정하면, 먼저 도달하는 옵션에 의해 마이크로 배치 크기가 제한돼요.
// Read a hard limit of 1 file per micro-batch
val df = spark.readStream
.format("iceberg")
.option("streaming-max-files-per-micro-batch", "1")
.load("database.table_name")
// Read files until the number of included rows >= 1000 per micro-batch
val df = spark.readStream
.format("iceberg")
.option("streaming-max-rows-per-micro-batch", "1000")
.load("database.table_name")
정보 (Info)
참고: 기본 트리거(즉 Trigger.ProcessingTime)를 사용하는 쿼리에서 마이크로 배치 크기를 제한하는 것 외에도, 비율 제한 옵션은 Trigger.AvailableNow를 사용하는 쿼리에 적용해서 사용 가능한 모든 소스 데이터의 일회성 처리를 더 나은 쿼리 확장성을 위해 여러 마이크로 배치로 나눌 수 있어요. 비율 제한 옵션은 더 이상 사용되지 않는(deprecated) Trigger.Once 트리거를 사용할 때는 무시돼요.
비동기 마이크로 배치 계획 (Asynchronous Micro-Batch Planning)
async-micro-batch-planning-enabled를 true로 설정하면 비동기 마이크로 배치 계획을 활성화할 수 있어요. 이 옵션을 활성화하면 아이스버그는 다음 마이크로 배치를 병렬로 계획하면서 현재 마이크로 배치를 처리하기 시작해요. 이렇게 하면 마이크로 배치 사이의 유휴 시간을 줄여 쿼리 처리량을 개선할 수 있어요. 더 높은 메모리 사용량과 증가된 스냅샷 감지 지연 같은 트레이드오프를 잘 따져봐야 해요.
사용자는 스파크 구성에서 확인할 수 있는 추가 옵션을 설정해서 비동기 마이크로 배치 계획의 동작을 제어할 수도 있어요.
스트리밍 쓰기 (Streaming Writes)
스트리밍 쿼리에서 아이스버그 테이블로 값을 쓰려면 DataStreamWriter를 사용해요.
data.writeStream
.format("iceberg")
.outputMode("append")
.trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
.option("checkpointLocation", checkpointPath)
.toTable("database.table_name")
디렉터리 기반 하둡 카탈로그의 경우:
data.writeStream
.format("iceberg")
.outputMode("append")
.trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
.option("path", "hdfs://nn:8020/path/to/table")
.option("checkpointLocation", checkpointPath)
.start()
아이스버그는 append와 complete 출력 모드를 지원해요.
- append: 매 마이크로 배치의 행을 테이블에 추가해요
- complete: 매 마이크로 배치마다 테이블 내용을 교체해요
스트리밍 쿼리를 시작하기 전에 테이블을 먼저 만들어 두었는지 확인해주세요. 아이스버그 테이블을 만드는 방법은 SQL create table 문서를 참고해주세요.
아이스버그는 실험적인 연속 처리(continuous processing)를 지원하지 않아요. 출력을 "커밋"할 인터페이스를 제공하지 않기 때문이에요.
파티션 테이블 (Partitioned table)
아이스버그는 데이터를 쓰기 전에 작업(task) 단위로 파티션별로 데이터를 정렬해야 해요. 스파크에서는 파티션된 테이블에 대해 작업들이 스파크 파티션으로 나뉘어져요. 배치 쿼리의 경우 명시적 정렬을 해서 요구사항을 충족하는 것이 권장되지만(여기 참조), repartition과 sort는 스트리밍 워크로드에서 무거운 연산으로 여겨지므로 이 방식은 추가 지연을 가져와요. 추가 지연을 피하려면 fanout 작성자를 활성화해서 그 요구사항을 제거할 수 있어요.
data.writeStream
.format("iceberg")
.outputMode("append")
.trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
.option("fanout-enabled", "true")
.option("checkpointLocation", checkpointPath)
.toTable("database.table_name")
fanout 작성자는 파티션 값별로 파일을 열고, 쓰기 작업이 끝날 때까지 그 파일들을 닫지 않아요. 배치 쓰기에서는 fanout 작성자 사용을 피해주세요. 배치 워크로드에서는 출력 행에 대한 명시적 정렬이 저렴하기 때문이에요.
스트리밍 테이블의 유지보수 (Maintenance for streaming tables)
스트리밍 쓰기는 새 테이블 버전을 빠르게 만들 수 있고, 그 버전들을 추적할 많은 테이블 메타데이터를 만들어내요. 커밋 비율을 조정하고, 오래된 스냅샷을 만료시키고, 메타데이터 파일을 자동으로 정리해서 메타데이터를 유지 관리하는 것을 적극 권장해요.
커밋 비율 조정 (Tune the rate of commits)
커밋 비율이 높으면 데이터 파일, 매니페스트, 스냅샷이 많이 생기고 추가 유지보수가 필요해져요. 트리거 간격을 최소 1분으로 두고, 필요하면 간격을 늘리는 것을 권장해요.
Structured Streaming Programming Guide의 triggers 섹션에서 간격을 구성하는 방법을 설명하고 있어요.
오래된 스냅샷 만료 (Expire old snapshots)
테이블에 쓰여진 각 배치는 새 스냅샷을 만들어요. 아이스버그는 스냅샷이 만료될 때까지 테이블 메타데이터에서 스냅샷을 추적해요. 잦은 커밋으로 스냅샷이 빠르게 누적되므로, 스트리밍 쿼리가 쓰는 테이블은 정기적으로 유지관리하는 것을 적극 권장해요. 스냅샷 만료는 더 이상 필요 없는 메타데이터와 데이터 파일을 제거하는 절차예요. 기본적으로 이 절차는 5일보다 오래된 스냅샷을 만료시켜요.
데이터 파일 컴팩션 (Compacting data files)
스트리밍 프로세스에서 쓰여지는 데이터 양은 보통 적어서, 테이블 메타데이터가 많은 작은 파일을 추적하게 될 수 있어요. 작은 파일을 더 큰 파일로 컴팩션하면 테이블이 필요로 하는 메타데이터가 줄어들고 쿼리 효율이 높아져요. 아이스버그와 스파크는 rewrite_data_files 프로시저를 제공해요.
매니페스트 재작성 (Rewrite manifests)
스트리밍 워크로드의 쓰기 지연을 최적화하기 위해, 아이스버그는 매니페스트를 자동으로 컴팩션하지 않는 "빠른" 추가로 새 스냅샷을 쓸 수 있어요. 이렇게 하면 작은 매니페스트 파일이 많이 생길 수 있어요. 아이스버그는 매니페스트 파일의 수를 재작성해서 쿼리 성능을 개선할 수 있어요. 아이스버그와 스파크는 rewrite_manifests 프로시저를 제공해요.