Apache Spark

配置

师成师成· 更新于 2026-09-28· 阅读 31 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Spark 配置

Catalog

Spark 提供了一套 API,可用于插入用于加载、创建和管理 Iceberg 表的 table catalog。Spark catalog 通过设置 spark.sql.catalog 下的 Spark 属性进行配置。

下面的配置会创建一个名为 hive_prod 的 Iceberg catalog,用于从 Hive metastore 加载表:

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_prod 的 REST 目录的示例,该目录从 REST URL http://localhost:8080 加载表:

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

Iceberg 还支持 HDFS 中基于目录的 catalog,可通过 type=hadoop 进行配置:

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

信息

基于 Hive 的目录只会加载 Iceberg 表。要在同一个 Hive 元存储中加载非 Iceberg 表,请使用会话目录。

目录配置

通过添加属性 spark.sql.catalog.(catalog-name) 并将其值设置为实现类,即可创建并命名一个目录。

Iceberg 提供两种实现:

  • org.apache.iceberg.spark.SparkCatalog 支持将 Hive Metastore 或 Hadoop 数据仓库作为目录
  • org.apache.iceberg.spark.SparkSessionCatalog 在 Spark 内置目录的基础上增加对 Iceberg 表的支持,并将非 Iceberg 表的请求委托给内置目录

两种目录都通过嵌套在目录名称下的属性进行配置。Hive 和 Hadoop 的常用配置属性如下:

属性值描述
spark.sql.catalog.catalog-name.typehive、hadoop、rest、glue、jdbc 或 nessie底层的 Iceberg catalog 实现:HiveCatalog、HadoopCatalog、RESTCatalog、GlueCatalog、JdbcCatalog、NessieCatalog;若使用自定义 catalog,则不设置该项
spark.sql.catalog.catalog-name.catalog-impl自定义的 Iceberg catalog 实现。如果 type 为 null,则 catalog-impl 不能为空。
spark.sql.catalog.catalog-name.io-impl自定义的 FileIO 实现。
spark.sql.catalog.catalog-name.metrics-reporter-impl自定义的 MetricsReporter 实现。
spark.sql.catalog.catalog-name.default-namespacedefaultcatalog 的默认当前命名空间
spark.sql.catalog.catalog-name.urithrift://host:porthive 类型 catalog 的 Hive metastore URL,REST 类型 catalog 的 REST URL
spark.sql.catalog.catalog-name.warehousehdfs://nn:8020/warehouse/pathwarehouse 目录的基础路径
spark.sql.catalog.catalog-name.cache-enabledtrue 或 false是否启用 catalog 缓存,默认值为 true
spark.sql.catalog.catalog-name.cache.expiration-interval-ms30000(30 秒)缓存的 catalog 条目经过该时长后过期;仅当 cache-enabled 为 true 时生效。-1 表示禁用缓存过期,0 表示完全禁用缓存(无论 cache-enabled 的取值如何)。默认值为 30000(30 秒)
spark.sql.catalog.catalog-name.table-default.propertyKey属性键 propertyKey 的默认 Iceberg 表属性值,若未被覆盖,将设置在由该 catalog 创建的表上
spark.sql.catalog.catalog-name.table-override.propertyKey属性键 propertyKey 的强制 Iceberg 表属性值,用户在创建表时无法覆盖
spark.sql.catalog.catalog-name.view-default.propertyKey属性键 propertyKey 的默认 Iceberg 视图属性值,若未被覆盖,将设置在由该 catalog 创建的视图上
spark.sql.catalog.catalog-name.view-override.propertyKey属性键 propertyKey 的强制 Iceberg 视图属性值,用户在创建视图时无法覆盖
spark.sql.catalog.catalog-name.use-nullable-query-schematrue 或 false使用 CTAS 和 RTAS 创建表时是否保留字段的可空性。设置为 true 时,所有字段都将被标记为可空;设置为 false 时,将保留字段的可空性。默认值为 true。该配置适用于 Spark 3.5 及以上版本。

其他属性可参见通用的目录配置。

使用目录

目录名称用于在 SQL 查询中标识表。在上面的示例中,hive_prod 和 hadoop_prod 可用作前缀,与数据库名和表名组合,从而从相应的目录中加载这些表。

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

Spark 3 会记录当前的 catalog 和命名空间,它们在表名中可以省略。

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

要查看当前目录和命名空间,请运行 SHOW CURRENT NAMESPACE。

替换会话目录

要为 Spark 内置目录添加 Iceberg 表支持,请将 spark_catalog 配置为使用 Iceberg 的 SparkSessionCatalog。

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

Spark 的内置 catalog 支持 Hive Metastore 中已有的 v1 和 v2 表。该配置使 Spark 使用 Iceberg 的 SparkSessionCatalog 作为该会话 catalog 的封装。当某个表不是 Iceberg 表时,将改用内置 catalog 来加载它。

通过此配置,Iceberg 表和非 Iceberg 表可以共用同一个 Hive Metastore。

当你希望 spark_catalog 在同一个 metastore 中同时处理 Iceberg 表和非 Iceberg 表时,SparkSessionCatalog 非常有用。

注意

Spark 4.2.0 之前的版本不支持会话 catalog 中的 V2Function。详情参见 SPARK-54760(apache/spark#53531)。因此,像 system.bucket、system.days 和 system.iceberg_version 这类基于 catalog 的 SQL 函数无法通过 spark_catalog 使用。要绕过此限制,请配置一个使用 org.apache.iceberg.spark.SparkCatalog 的独立 Iceberg catalog,并通过该 catalog 调用这些函数。

使用特定 catalog 的 Hadoop 配置项

与通过 spark.hadoop.* 配置 Hadoop 属性类似,在使用 Spark 时,可以通过添加带有前缀 spark.sql.catalog.(catalog-name).hadoop.* 的属性来设置每个 catalog 独立的 Hadoop 配置项。这些属性的优先级高于通过 spark.hadoop.* 全局配置的值,并且只会影响 Iceberg 表。

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

加载自定义目录

Spark 支持通过指定 catalog-impl 属性来加载自定义的 Iceberg Catalog 实现。示例如下:

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 扩展

Iceberg 0.11.0 及更高版本为 Spark 添加了一个扩展模块,用于引入新的 SQL 命令,例如用于存储过程的 CALL 以及 ALTER TABLE ... WRITE ORDERED BY。

使用这些 SQL 命令需要通过以下 Spark 属性将 Iceberg 扩展添加到你的 Spark 环境中:

Spark 扩展属性Iceberg 扩展实现
spark.sql.extensionsorg.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

运行时配置

配置项的优先级

Iceberg 允许在不同层级指定配置。读写操作的生效配置依据以下优先级顺序确定:

  1. DataSource API 读/写选项 – 在读写操作中通过 .option(...) 显式传递。
  2. Spark Session 配置 – 通过 spark.conf.set(...)、spark-defaults.conf 或 spark-submit 中的 --conf 在 Spark 中全局设置。
  3. 表属性 – 通过 ALTER TABLE SET TBLPROPERTIES 定义在 Iceberg 表上。
  4. 默认值。

如果某个设置在更高层级未定义,则会回退到下一层级。这样既保持了灵活性,又能在需要时启用全局默认值。

Spark SQL 选项

Iceberg 支持使用 Spark SQL 配置选项来设置各种全局行为。这些选项可以通过 spark.conf、SparkSession 设置或 Spark 提交参数进行配置。例如:

// 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 选项默认值说明
spark.sql.iceberg.vectorization.enabled表默认值启用数据文件的向量化读取
spark.sql.iceberg.check-nullabilitytrue校验写入 schema 的可空性是否与表的可空性一致
spark.sql.iceberg.check-orderingtrue校验写入 schema 的列顺序是否与表 schema 的顺序一致
spark.sql.iceberg.planning.preserve-data-groupingfalse为 true 时,将同一分区的扫描任务集中放置在同一个读取分片中,用于存储分区连接(Storage Partitioned Joins)
spark.sql.iceberg.aggregate-push-down.enabledtrue启用聚合函数(MAX、MIN、COUNT)下推
spark.sql.iceberg.distribution-mode参见 Spark 写入控制写入过程中的数据分发策略
spark.wap.idnullWrite-Audit-Publish 快照暂存 ID
spark.wap.branchnull用于快照提交的 WAP 分支名称
spark.sql.iceberg.shred-variants表默认值为 true 时,变体(variant)列采用拆分(shredded)Parquet 编码写入,以提升查询性能
spark.sql.iceberg.variant-inference-buffer-size表默认值启用变体拆分时,用于 schema 推断的缓冲行数
spark.sql.iceberg.compression-codec表默认值写入压缩编解码器(例如 zstd、snappy)
spark.sql.iceberg.compression-level表默认值Parquet/Avro 的压缩级别
spark.sql.iceberg.compression-strategy表默认值ORC 的压缩策略
spark.sql.iceberg.data-planning-modeAUTO数据文件的扫描规划模式(AUTO、LOCAL、DISTRIBUTED)
spark.sql.iceberg.delete-planning-modeAUTO删除文件的扫描规划模式(AUTO、LOCAL、DISTRIBUTED)
spark.sql.iceberg.advisory-partition-size表默认值启用 Spark 自适应查询执行(Adaptive Query Execution)时,写入表所使用的建议大小(字节),用于确定输出文件大小
spark.sql.iceberg.locality.enabledfalse向 Spark 报告执行器(executor)上的任务放置位置信息
spark.sql.iceberg.executor-cache.enabledtrue启用执行器端缓存(当前用于缓存删除文件 Delete Files)
spark.sql.iceberg.executor-cache.timeout10执行器缓存条目的超时时间(分钟)
spark.sql.iceberg.executor-cache.max-entry-size67108864 (64MB)单个缓存条目的最大大小(字节)
spark.sql.iceberg.executor-cache.max-total-size134217728 (128MB)执行器缓存的最大总大小(字节)
spark.sql.iceberg.executor-cache.locality.enabledfalse启用感知数据局部性的执行器缓存
spark.sql.iceberg.merge-schemafalse启用修改表 schema 以匹配写入 schema 的能力,仅会添加缺失的列
spark.sql.iceberg.report-column-statstrue在可用时,将 Puffin 表统计信息提供给 Spark 的基于成本的优化器(CBO);需要启用 CBO 才能生效
spark.sql.iceberg.async-micro-batch-planning-enabledfalse启用异步微批规划,通过预取文件扫描任务来降低规划延迟

读取选项

Spark 读取选项在配置 DataFrameReader 时传入,如下所示:

// time travel
spark.read
    .option("snapshot-id", 10963874102873L)
    .table("catalog.db.table")
Spark 选项默认值说明
snapshot-id(最新)要读取的表快照的快照 ID
as-of-timestamp(最新)以毫秒为单位的时间戳;将使用该时刻当前存在的快照
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(无)要开始流式处理的毫秒时间戳;如果早于最旧的已知祖先快照,则使用最旧的快照
streaming-max-files-per-micro-batchINT_MAX每个微批处理的最大文件数
streaming-max-rows-per-micro-batchINT_MAX每个微批处理的“软上限”行数;始终包含下一个未处理文件中的所有行,若加入其他文件会导致超出软上限则将其排除
async-micro-batch-planning-enabledfalse启用异步微批处理规划,通过预取文件扫描任务来降低规划延迟
streaming-snapshot-polling-interval-ms30000覆盖异步规划器刷新和检测新快照的轮询时间间隔。仅在设置了 async-micro-batch-planning-enabled 时生效
async-queue-preload-file-limit100覆盖初始加载到后台队列的文件数量。可通过调优避免队列饥饿。仅在设置了 async-micro-batch-planning-enabled 时生效
async-queue-preload-row-limit100000覆盖初始加载到后台队列的行数。可通过调优避免队列饥饿。仅在设置了 async-micro-batch-planning-enabled 时生效

写入选项

Spark 写入选项在配置 DataFrameWriterV2 时传入,示例如下:

// write with Avro instead of Parquet
df.writeTo("catalog.db.table")
    .option("write-format", "avro")
    .option("snapshot-property.key", "value")
    .append()
Spark 选项默认值说明
write-formatTable write.format.default本次写入操作使用的文件格式:parquet、avro 或 orc
target-file-size-bytes按表属性而定覆盖该表的 write.target-file-size-bytes
check-nullabilitytrue对字段启用可空性检查
snapshot-property.custom-keynull在快照摘要中添加 custom-key 及其对应值的条目(snapshot-property. 前缀仅在 DSv2 中需要)
fanout-enabledfalse覆盖该表的 write.spark.fanout.enabled
check-orderingtrue检查输入 schema 与表 schema 是否相同
isolation-levelnullDataframe 覆写操作所需的隔离级别。null => 不做检查(适用于幂等写入),serializable => 检查目标分区中是否存在并发插入或删除,snapshot => 检查目标分区中是否存在并发删除。
validate-from-snapshot-idnull若设置了隔离级别,则指定用于检查写入该表的并发写冲突的基础快照 ID。该快照应是在对该表执行任何读取操作之前的快照。可通过 Table API 或 Snapshots 表 获取。若为 null,则使用该表已知的最早快照。
compression-codecTable write.(fileformat).compression-codec覆盖本次写入使用的该表压缩编解码器
compression-levelTable write.(fileformat).compression-level覆盖本次写入中 Parquet 和 Avro 表的压缩级别
compression-strategyTable write.orc.compression-strategy覆盖本次写入中 ORC 表的压缩策略
distribution-mode默认值参见 Spark Writes覆盖本次写入使用的该表分布模式
delete-granularityfile覆盖本次写入使用的该表删除粒度
shred-variantsfalse覆盖本次写入使用的该表 write.parquet.shred-variants
variant-inference-buffer-size100覆盖本次写入使用的该表 write.parquet.variant-inference-buffer-size

CommitMetadata 提供了一个接口,可在 SQL 执行过程中向快照摘要(snapshot summary)中添加自定义元数据,这对于审计或变更跟踪等用途非常有价值。如果属性名以 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);

评论

登录后参与评论

正在加载评论…