配置
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:8080Iceberg 还支持 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.type | hive、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-namespace | default | catalog 的默认当前命名空间 |
| spark.sql.catalog.catalog-name.uri | thrift://host:port | hive 类型 catalog 的 Hive metastore URL,REST 类型 catalog 的 REST URL |
| spark.sql.catalog.catalog-name.warehouse | hdfs://nn:8020/warehouse/path | warehouse 目录的基础路径 |
| spark.sql.catalog.catalog-name.cache-enabled | true 或 false | 是否启用 catalog 缓存,默认值为 true |
| spark.sql.catalog.catalog-name.cache.expiration-interval-ms | 30000(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-schema | true 或 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_prodSpark 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 = hiveSpark 的内置 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-valueSQL 扩展
Iceberg 0.11.0 及更高版本为 Spark 添加了一个扩展模块,用于引入新的 SQL 命令,例如用于存储过程的 CALL 以及 ALTER TABLE ... WRITE ORDERED BY。
使用这些 SQL 命令需要通过以下 Spark 属性将 Iceberg 扩展添加到你的 Spark 环境中:
| Spark 扩展属性 | Iceberg 扩展实现 |
|---|---|
spark.sql.extensions | org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions |
运行时配置
配置项的优先级
Iceberg 允许在不同层级指定配置。读写操作的生效配置依据以下优先级顺序确定:
- DataSource API 读/写选项 – 在读写操作中通过
.option(...)显式传递。 - Spark Session 配置 – 通过
spark.conf.set(...)、spark-defaults.conf或 spark-submit 中的--conf在 Spark 中全局设置。 - 表属性 – 通过
ALTER TABLE SET TBLPROPERTIES定义在 Iceberg 表上。 - 默认值。
如果某个设置在更高层级未定义,则会回退到下一层级。这样既保持了灵活性,又能在需要时启用全局默认值。
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-nullability | true | 校验写入 schema 的可空性是否与表的可空性一致 |
| spark.sql.iceberg.check-ordering | true | 校验写入 schema 的列顺序是否与表 schema 的顺序一致 |
| 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 | 参见 Spark 写入 | 控制写入过程中的数据分发策略 |
| spark.wap.id | null | Write-Audit-Publish 快照暂存 ID |
| spark.wap.branch | null | 用于快照提交的 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-mode | AUTO | 数据文件的扫描规划模式(AUTO、LOCAL、DISTRIBUTED) |
| spark.sql.iceberg.delete-planning-mode | AUTO | 删除文件的扫描规划模式(AUTO、LOCAL、DISTRIBUTED) |
| spark.sql.iceberg.advisory-partition-size | 表默认值 | 启用 Spark 自适应查询执行(Adaptive Query Execution)时,写入表所使用的建议大小(字节),用于确定输出文件大小 |
| spark.sql.iceberg.locality.enabled | false | 向 Spark 报告执行器(executor)上的任务放置位置信息 |
| 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 | 启用修改表 schema 以匹配写入 schema 的能力,仅会添加缺失的列 |
| spark.sql.iceberg.report-column-stats | true | 在可用时,将 Puffin 表统计信息提供给 Spark 的基于成本的优化器(CBO);需要启用 CBO 才能生效 |
| spark.sql.iceberg.async-micro-batch-planning-enabled | false | 启用异步微批规划,通过预取文件扫描任务来降低规划延迟 |
读取选项
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-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 时生效 |
写入选项
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-format | Table write.format.default | 本次写入操作使用的文件格式:parquet、avro 或 orc |
| target-file-size-bytes | 按表属性而定 | 覆盖该表的 write.target-file-size-bytes |
| check-nullability | true | 对字段启用可空性检查 |
| snapshot-property.custom-key | null | 在快照摘要中添加 custom-key 及其对应值的条目(snapshot-property. 前缀仅在 DSv2 中需要) |
| fanout-enabled | false | 覆盖该表的 write.spark.fanout.enabled |
| check-ordering | true | 检查输入 schema 与表 schema 是否相同 |
| isolation-level | null | Dataframe 覆写操作所需的隔离级别。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 | 默认值参见 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 执行过程中向快照摘要(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);评论
登录后参与评论
KnowForge