스파크 구성

스파크 구성 (Configuration)

이 문서에서는 스파크에서 아이스버그 카탈로그를 구성하고 런타임 옵션을 설정하는 방법을 알려드릴게요. spark.sql.catalog 아래에 카탈로그를 플러그인하는 방법, SQL 확장, 읽기·쓰기 옵션, 그리고 설정 우선순위까지 표와 예시를 통해 자세히 살펴볼게요.

출처: 문서

본문

카탈로그 (Catalogs)

스파크는 아이스버그 테이블을 로드, 생성, 관리하는 데 사용되는 테이블 카탈로그를 플러그인하는 API를 추가해요. 스파크 카탈로그는 spark.sql.catalog 아래에 스파크 속성을 설정해서 구성돼요.

이것은 Hive metastore에서 테이블을 로드하는 hive_prod라는 아이스버그 카탈로그를 만들어요.

spark.sql.catalog.hive_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type = hive
spark.sql.catalog.hive_prod.uri = thrift://metastore-host:port
# omit uri to use the same URI as Spark: hive.metastore.uris in hive-site.xml

아래는 REST URL http://localhost:8080에서 테이블을 로드하는 rest_prod라는 REST 카탈로그의 예시예요.

spark.sql.catalog.rest_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.rest_prod.type = rest
spark.sql.catalog.rest_prod.uri = http://localhost:8080

아이스버그는 type=hadoop으로 구성할 수 있는 HDFS의 디렉터리 기반 카탈로그도 지원해요.

spark.sql.catalog.hadoop_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hadoop_prod.type = hadoop
spark.sql.catalog.hadoop_prod.warehouse = hdfs://nn:8020/warehouse/path

정보 (Info)

Hive 기반 카탈로그는 아이스버그 테이블만 로드해요. 같은 Hive metastore에서 non-Iceberg 테이블을 로드하려면 세션 카탈로그(session catalog)를 사용해주세요.

카탈로그 구성 (Catalog configuration)

카탈로그는 값으로 구현 클래스를 가진 spark.sql.catalog.(catalog-name) 속성을 추가해서 만들어지고 이름이 붙여져요.

아이스버그는 두 가지 구현을 제공해요.

  • org.apache.iceberg.spark.SparkCatalog: Hive Metastore 또는 Hadoop warehouse를 카탈로그로 지원해요
  • org.apache.iceberg.spark.SparkSessionCatalog: 스파크 내장 카탈로그에 아이스버그 테이블 지원을 추가하고, non-Iceberg 테이블은 내장 카탈로그에 위임해요

두 카탈로그 모두 카탈로그 이름 아래에 중첩된 속성으로 구성돼요. Hive와 Hadoop의 공통 구성 속성은 다음과 같아요.

속성 설명
spark.sql.catalog.catalog-name.type hive, hadoop, rest, glue, jdbc 또는 nessie 기본 아이스버그 카탈로그 구현, HiveCatalog, HadoopCatalog, RESTCatalog, GlueCatalog, JdbcCatalog, NessieCatalog 또는 커스텀 카탈로그를 사용할 때는 미설정
spark.sql.catalog.catalog-name.catalog-impl 커스텀 아이스버그 카탈로그 구현. type이 null이면 catalog-impl은 null이 아니어야 해요.
spark.sql.catalog.catalog-name.io-impl 커스텀 FileIO 구현.
spark.sql.catalog.catalog-name.metrics-reporter-impl 커스텀 MetricsReporter 구현.
spark.sql.catalog.catalog-name.default-namespace default 카탈로그의 기본 현재 네임스페이스
spark.sql.catalog.catalog-name.uri thrift://host:port hive 타입 카탈로그의 Hive metastore URL, REST 타입 카탈로그의 REST URL
spark.sql.catalog.catalog-name.warehouse hdfs://nn:8020/warehouse/path warehouse 디렉터리의 기본 경로
spark.sql.catalog.catalog-name.cache-enabled true 또는 false 카탈로그 캐시 활성화 여부, 기본값은 true
spark.sql.catalog.catalog-name.cache.expiration-interval-ms 30000 (30 seconds) 캐시된 카탈로그 항목이 만료되는 시간. cache-enabled가 true일 때만 유효해요. -1은 캐시 만료를 비활성화하고, 0은 cache-enabled와 무관하게 캐시를 완전히 비활성화해요. 기본값은 30000 (30 seconds)
spark.sql.catalog.catalog-name.table-default.propertyKey 속성 키 propertyKey에 대한 기본 아이스버그 테이블 속성 값. 이 카탈로그로 생성된 테이블에 재정의되지 않으면 설정돼요
spark.sql.catalog.catalog-name.table-override.propertyKey 속성 키 propertyKey에 대한 강제 아이스버그 테이블 속성 값. 사용자가 테이블 생성 시 재정의할 수 없어요
spark.sql.catalog.catalog-name.view-default.propertyKey 속성 키 propertyKey에 대한 기본 아이스버그 뷰 속성 값. 이 카탈로그로 생성된 뷰에 재정의되지 않으면 설정돼요
spark.sql.catalog.catalog-name.view-override.propertyKey 속성 키 propertyKey에 대한 강제 아이스버그 뷰 속성 값. 사용자가 뷰 생성 시 재정의할 수 없어요
spark.sql.catalog.catalog-name.use-nullable-query-schema true 또는 false CTAS와 RTAS로 테이블을 만들 때 필드의 null 허용 여부를 보존할지 여부. true로 설정하면 모든 필드가 nullable로 표시돼요. false로 설정하면 필드의 null 허용 여부가 보존돼요. 기본값은 true. Spark 3.5 이상에서 사용 가능

추가 속성은 공통 카탈로그 구성(common catalog configuration)에서 찾을 수 있어요.

카탈로그 사용 (Using catalogs)

카탈로그 이름은 테이블을 식별하기 위해 SQL 쿼리에서 사용돼요. 위 예시에서 hive_prod와 hadoop_prod는 그 카탈로그에서 로드될 데이터베이스와 테이블 이름에 프리픽스를 붙이는 데 사용할 수 있어요.

SELECT * FROM hive_prod.db.table; -- load db.table from catalog hive_prod

Spark 3은 현재 카탈로그와 네임스페이스를 추적하고, 이는 테이블 이름에서 생략할 수 있어요.

USE hive_prod.db;
SELECT * FROM table; -- load db.table from catalog hive_prod

현재 카탈로그와 네임스페이스를 보려면 SHOW CURRENT NAMESPACE를 실행해요.

세션 카탈로그 교체 (Replacing the session catalog)

스파크 내장 카탈로그에 아이스버그 테이블 지원을 추가하려면 spark_catalog을 아이스버그의 SparkSessionCatalog으로 구성해요.

spark.sql.catalog.spark_catalog = org.apache.iceberg.spark.SparkSessionCatalog
spark.sql.catalog.spark_catalog.type = hive

스파크 내장 카탈로그는 Hive Metastore에서 추적되는 기존 v1과 v2 테이블을 지원해요. 이것은 스파크가 그 세션 카탈로그를 감싸는 래퍼로 아이스버그의 SparkSessionCatalog를 사용하도록 구성해요. 테이블이 아이스버그 테이블이 아니면 내장 카탈로그를 사용해서 대신 로드돼요.

이 구성은 같은 Hive Metastore를 아이스버그와 non-Iceberg 테이블 모두에 사용할 수 있어요.

SparkSessionCatalog는 같은 metastore에서 아이스버그와 non-Iceberg 테이블 모두로 spark_catalog을 작업하고 싶을 때 유용해요.

참고 (Note)

4.2.0 이전의 스파크는 세션 카탈로그에서 V2Function을 지원하지 않아요. 자세한 내용은 SPARK-54760 (apache/spark#53531)을 참고해주세요. 결과적으로 system.bucket, system.days, system.iceberg_version 같은 카탈로그 범위 SQL 함수는 spark_catalog을 통해 사용할 수 없어요. 이 제한을 우회하려면 org.apache.iceberg.spark.SparkCatalog로 별도의 아이스버그 카탈로그를 구성하고 그 카탈로그를 통해 호출해주세요.

카탈로그별 Hadoop 구성 값 사용 (Using catalog specific Hadoop configuration values)

spark.hadoop.을 사용해서 Hadoop 속성을 구성하는 것과 비슷하게, 카탈로그에 대한 Hadoop 구성 값을 설정하려면 spark.sql.catalog.(catalog-name).hadoop. 프리픽스로 카탈로그 속성을 추가해서 사용할 수 있어요. 이 속성들은 spark.hadoop.*로 전역적으로 구성된 값보다 우선하고, 아이스버그 테이블에만 영향을 줘요.

spark.sql.catalog.hadoop_prod.hadoop.fs.s3a.endpoint = http://aws-local:9000

커스텀 카탈로그 로드 (Loading a custom catalog)

스파크는 catalog-impl 속성을 지정해서 커스텀 아이스버그 카탈로그 구현을 로드하는 것을 지원해요. 예시는 다음과 같아요.

spark.sql.catalog.custom_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.custom_prod.catalog-impl = com.my.custom.CatalogImpl
spark.sql.catalog.custom_prod.my-additional-catalog-config = my-value

SQL 확장 (SQL Extensions)

Iceberg 0.11.0 이상은 스파크에 확장 모듈을 추가해서 저장 프로시저용 CALL이나 ALTER TABLE ... WRITE ORDERED BY 같은 새 SQL 명령을 제공해요.

그런 SQL 명령을 사용하려면 다음 스파크 속성을 사용해서 스파크 환경에 아이스버그 확장을 추가해야 해요.

스파크 확장 속성 아이스버그 확장 구현
spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

런타임 구성 (Runtime configuration)

구성 설정의 우선순위 (Precedence of Configuration Settings)

아이스버그는 서로 다른 수준에서 구성을 지정할 수 있게 해줘요. 읽기나 쓰기 연산의 유효 구성은 다음 우선순위 순서에 따라 결정돼요.

  1. DataSource API 읽기/쓰기 옵션 – 읽기/쓰기 연산에서 .option(...)에 명시적으로 전달된 것.
  2. 스파크 세션 구성 – spark.conf.set(...), spark-defaults.conf, 또는 spark-submit의 --conf로 스파크에 전역 설정된 것.
  3. 테이블 속성 – ALTER TABLE SET TBLPROPERTIES로 아이스버그 테이블에 정의된 것.
  4. 기본값.

설정이 더 높은 수준에 정의되지 않으면 다음 수준이 폴백으로 사용돼요. 이는 필요할 때 전역 기본값을 활성화하면서도 유연성을 허용해요.

스파크 SQL 옵션 (Spark SQL Options)

아이스버그는 다양한 전역 동작을 스파크 SQL 구성 옵션으로 설정하는 것을 지원해요. 이것들은 spark.conf, SparkSession 설정, 스파크 제출 인자로 설정할 수 있어요. 예를 들어:

// disabling vectorization
val spark = SparkSession.builder()
  .appName("IcebergExample")
  .master("local[*]")
  .config("spark.sql.catalog.my_catalog", "org.apache.iceberg.spark.SparkCatalog")
  .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
  .config("spark.sql.iceberg.vectorization.enabled", "false")
  .getOrCreate()
스파크 옵션 기본값 설명
spark.sql.iceberg.vectorization.enabled Table default 데이터 파일의 벡터화 읽기 활성화
spark.sql.iceberg.check-nullability true 쓰기 스키마의 null 허용 여부가 테이블의 null 허용 여부와 일치하는지 검증
spark.sql.iceberg.check-ordering true 쓰기 스키마 컬럼 순서가 테이블 스키마 순서와 일치하는지 검증
spark.sql.iceberg.planning.preserve-data-grouping false true이면 같은 파티션의 스캔 작업을 같은 읽기 스플릿에 공동 배치. Storage Partitioned Joins에 사용
spark.sql.iceberg.aggregate-push-down.enabled true 집계 함수(MAX, MIN, COUNT)의 푸시다운 활성화
spark.sql.iceberg.distribution-mode See Spark Writes 쓰기 중 분포 전략 제어
spark.wap.id null Write-Audit-Publish 스냅샷 스테이징 ID
spark.wap.branch null 스냅샷 커밋을 위한 WAP 브랜치 이름
spark.sql.iceberg.shred-variants Table default true이면 variant 컬럼이 쿼리 성능 향상을 위해 분쇄된(shredded) Parquet 인코딩으로 쓰여져요
spark.sql.iceberg.variant-inference-buffer-size Table default variant 분쇄가 활성화될 때 스키마 추론을 위해 버퍼링할 행 수
spark.sql.iceberg.compression-codec Table default 쓰기 압축 코덱 (예: zstd, snappy)
spark.sql.iceberg.compression-level Table default Parquet/Avro의 압축 수준
spark.sql.iceberg.compression-strategy Table default ORC의 압축 전략
spark.sql.iceberg.data-planning-mode AUTO 데이터 파일의 스캔 계획 모드 (AUTO, LOCAL, DISTRIBUTED)
spark.sql.iceberg.delete-planning-mode AUTO 삭제 파일의 스캔 계획 모드 (AUTO, LOCAL, DISTRIBUTED)
spark.sql.iceberg.advisory-partition-size Table default 스파크의 Adaptive Query Execution이 활성화될 때 테이블 쓰기에 사용되는 권고 크기(바이트). 출력 파일 크기를 정하는 데 사용
spark.sql.iceberg.locality.enabled false 실행자에 대한 스파크 작업 배치의 지역성(locality) 정보 보고
spark.sql.iceberg.executor-cache.enabled true 실행자 측 캐시 활성화 (현재 Delete Files 캐싱에 사용)
spark.sql.iceberg.executor-cache.timeout 10 실행자 캐시 항목의 타임아웃(분)
spark.sql.iceberg.executor-cache.max-entry-size 67108864 (64MB) 캐시 항목당 최대 크기(바이트)
spark.sql.iceberg.executor-cache.max-total-size 134217728 (128MB) 실행자 캐시의 최대 총 크기(바이트)
spark.sql.iceberg.executor-cache.locality.enabled false 지역성 인식 실행자 캐시 사용 활성화
spark.sql.iceberg.merge-schema false 쓰기 스키마와 일치하도록 테이블 스키마 수정 활성화. 누락된 컬럼만 추가해요
spark.sql.iceberg.report-column-stats true Puffin 테이블 통계를 사용할 수 있으면 스파크의 Cost Based Optimizer에 보고. CBO가 활성화돼야 효과가 있어요
spark.sql.iceberg.async-micro-batch-planning-enabled false 파일 스캔 작업을 미리 가져와 계획 지연을 줄이는 비동기 마이크로 배치 계획 활성화

읽기 옵션 (Read options)

스파크 읽기 옵션은 DataFrameReader를 구성할 때 전달돼요.

// time travel
spark.read
    .option("snapshot-id", 10963874102873L)
    .table("catalog.db.table")
스파크 옵션 기본값 설명
snapshot-id (latest) 읽을 테이블 스냅샷의 스냅샷 ID
as-of-timestamp (latest) 밀리초 단위 타임스탬프; 사용되는 스냅샷은 이 시점의 현재 스냅샷
split-size 테이블 속성 기준 이 테이블의 read.split.target-size와 read.split.metadata-target-size를 재정의
lookback 테이블 속성 기준 이 테이블의 read.split.planning-lookback을 재정의
file-open-cost 테이블 속성 기준 이 테이블의 read.split.open-file-cost를 재정의
vectorization-enabled 테이블 속성 기준 이 테이블의 read.parquet.vectorization.enabled를 재정의
batch-size 테이블 속성 기준 이 테이블의 read.parquet.vectorization.batch-size를 재정의
stream-from-timestamp (none) 스트리밍할 밀리초 단위 타임스탬프; 알려진 가장 오래된 조상 스냅샷보다 이전이면 가장 오래된 것이 사용됨
streaming-max-files-per-micro-batch INT_MAX 마이크로 배치당 최대 파일 수
streaming-max-rows-per-micro-batch INT_MAX 마이크로 배치당 "소프트 최대" 행 수; 항상 다음 미처리 파일의 모든 행을 포함하고, 포함 시 소프트 최대를 초과하면 추가 파일을 제외함
async-micro-batch-planning-enabled false 파일 스캔 작업을 미리 가져와 계획 지연을 줄이는 비동기 마이크로 배치 계획 활성화
streaming-snapshot-polling-interval-ms 30000 비동기 플래너가 새 스냅샷을 갱신하고 감지하는 폴링 시간을 재정의. async-micro-batch-planning-enabled가 설정됐을 때만 영향
async-queue-preload-file-limit 100 초기에 백그라운드 큐에 로드되는 파일 수를 재정의. 큐 고갈을 방지하도록 조정. async-micro-batch-planning-enabled가 설정됐을 때만 영향
async-queue-preload-row-limit 100000 초기에 백그라운드 큐에 로드되는 행 수를 재정의. 큐 고갈을 방지하도록 조정. async-micro-batch-planning-enabled가 설정됐을 때만 영향

쓰기 옵션 (Write options)

스파크 쓰기 옵션은 DataFrameWriterV2를 구성할 때 전달돼요.

// write with Avro instead of Parquet
df.writeTo("catalog.db.table")
    .option("write-format", "avro")
    .option("snapshot-property.key", "value")
    .append()
스파크 옵션 기본값 설명
write-format Table write.format.default 이 쓰기 연산에 사용할 파일 포맷; parquet, avro 또는 orc
target-file-size-bytes 테이블 속성 기준 이 테이블의 write.target-file-size-bytes를 재정의
check-nullability true 필드에 nullable 검사 설정
snapshot-property.custom-key null 스냅샷 요약에 custom-key와 해당 값을 항목으로 추가 (snapshot-property. 프리픽스는 DSv2에서만 필요)
fanout-enabled false 이 테이블의 write.spark.fanout.enabled를 재정의
check-ordering true 입력 스키마와 테이블 스키마가 같은지 확인
isolation-level null Dataframe overwrite 연산에 대한 원하는 격리 수준. null => 검사 없음 (멱등 쓰기용), serializable => 대상 파티션의 동시 삽입이나 삭제 검사, snapshot => 대상 파티션의 동시 삭제 검사
validate-from-snapshot-id null 격리 수준이 설정되면, 테이블로의 동시 쓰기 충돌을 검사할 기준 스냅샷의 ID. 테이블에서 어떤 읽기보다 이전의 스냅샷이어야 해요. Table API나 Snapshots 테이블로 얻을 수 있어요. null이면 테이블의 가장 오래된 알려진 스냅샷이 사용됨
compression-codec Table write.(fileformat).compression-codec 이 쓰기의 테이블 압축 코덱을 재정의
compression-level Table write.(fileformat).compression-level Parquet 및 Avro 테이블의 이 쓰기 압축 수준을 재정의
compression-strategy Table write.orc.compression-strategy ORC 테이블의 이 쓰기 압축 전략을 재정의
distribution-mode defaults는 Spark Writes 참조 이 쓰기의 테이블 분포 모드를 재정의
delete-granularity file 이 쓰기의 테이블 삭제 세분성을 재정의
shred-variants false 이 쓰기의 write.parquet.shred-variants를 재정의
variant-inference-buffer-size 100 이 쓰기의 write.parquet.variant-inference-buffer-size를 재정의

CommitMetadata는 SQL 실행 중 스냅샷 요약에 커스텀 메타데이터를 추가하는 인터페이스를 제공해요. 이는 감사(auditing)나 변경 추적 같은 목적에 유용해요. 속성이 snapshot-property.로 시작하면 그 프리픽스가 각 속성에서 제거돼요. 예시는 다음과 같아요.

import org.apache.iceberg.spark.CommitMetadata;

Map<String, String> properties = Maps.newHashMap();
properties.put("property_key", "property_value");
CommitMetadata.withCommitProperties(properties,
        () -> {
            spark.sql("DELETE FROM " + tableName + " where id = 1");
            return 0;
        },
        RuntimeException.class);

더 알아보기 (Learn more)