스파크 구성
스파크 구성 (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)
아이스버그는 서로 다른 수준에서 구성을 지정할 수 있게 해줘요. 읽기나 쓰기 연산의 유효 구성은 다음 우선순위 순서에 따라 결정돼요.
- DataSource API 읽기/쓰기 옵션 – 읽기/쓰기 연산에서 .option(...)에 명시적으로 전달된 것.
- 스파크 세션 구성 – spark.conf.set(...), spark-defaults.conf, 또는 spark-submit의 --conf로 스파크에 전역 설정된 것.
- 테이블 속성 – ALTER TABLE SET TBLPROPERTIES로 아이스버그 테이블에 정의된 것.
- 기본값.
설정이 더 높은 수준에 정의되지 않으면 다음 수준이 폴백으로 사용돼요. 이는 필요할 때 전역 기본값을 활성화하면서도 유연성을 허용해요.
스파크 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);