为何写入Iceberg表的流处理变更需新Spark Session才能可见?
问题
使用Spark Scala以批处理模式创建Iceberg表后,执行含Merge Into操作的流处理写入,原批处理所用的Spark Session无法查看新变更,必须新建Spark Session才能看到表的实际状态。
批处理写入代码:df.writeTo(tableName).createOrReplace()
流处理写入相关方法:
override def write(df: DataFrame): StreamingQuery = { implicit val spark: SparkSession = df.sparkSession df.writeStream .format("iceberg") .trigger(Trigger.Once()) .option("fanout-enabled", "true") .option("checkpointLocation", checkpointLocation) .foreachBatch(mergeDFIntoTable _) .outputMode("update") .start() }
private def mergeDFIntoTable(df: DataFrame, batchId: Long): Unit = { df.createOrReplaceTempView("source_table") val mergeSQL: String = s""" |MERGE INTO $TableIdentifier t |USING source_table s |ON s.identifier = t.identifier |WHEN MATCHED THEN UPDATE SET * |WHEN NOT MATCHED THEN INSERT * |""".stripMargin logger.info(s"BATCH ID: $batchId | SQL: $mergeSQL") // The streaming query uses a cloned Spark Session from the original session that created it. df.sparkSession.sql(mergeSQL) df.unpersist() }
现象:使用原Spark Session执行spark.table(tableName).show(),批处理后和流处理后的结果一致,无法看到新增或更新的行;但新建Sessionval newSparkSession = spark.newSession()后执行newSparkSession.table(tableName).show(),就能看到流处理的变更。
原因分析
- Spark Session的元数据缓存机制:原Spark Session首次读取Iceberg表时,会将表的元数据(分区信息、快照版本等)缓存到Session级别的缓存中。流处理的Merge Into操作会生成Iceberg新快照,但原Session不会主动刷新缓存,查询时仍基于旧元数据返回结果。
- Iceberg的快照隔离特性:Iceberg通过快照(Snapshot)管理表版本,每个写操作都会生成新快照。原Session初始化时绑定了表的旧快照,除非主动触发元数据刷新,否则会一直使用该快照的数据视图。新建Session会重新读取表的最新元数据和快照,因此能看到最新变更。
- 流处理克隆Session的独立性:流处理
foreachBatch使用的是原Session的克隆Session,这个克隆Session的元数据变更不会同步回原Session。Merge Into操作在克隆Session中完成后,原Session完全感知不到元数据更新,查询结果自然还是旧数据。
内容的提问来源于stack exchange,提问作者Emilio
相关产品推荐
相关产品推荐

