You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

请说明代码中microBatchDF._jdf.sparkSession().sql的功能

Spark Streaming 实现Bronze到Silver表的Upsert及核心代码解析

这段代码的作用是通过Spark Streaming读取Bronze层表数据,以mrn作为唯一标识,对Silver层表执行Upsert(更新+插入)操作:当Silver表中存在相同mrn的记录且字段有变更时更新,不存在则插入新记录。

核心代码microBatchDF._jdf.sparkSession().sql的功能

这个调用的本质是通过Python API调用底层Spark的SQL执行引擎,执行Delta Lake的MERGE INTO语句完成Upsert,具体拆解:

  • microBatchDF:Spark Streaming微批处理中当前批次的待处理DataFrame
  • _jdf:获取Python DataFrame对应的底层Java DataFrame对象(Spark核心基于Java/Scala实现,Python API是封装层)
  • .sparkSession():获取当前DataFrame所属的SparkSession实例(即当前Spark执行环境的会话)
  • .sql():调用SparkSession的SQL执行能力,执行传入的MERGE INTO语句,完成对Silver表的更新和插入操作

这段SQL逻辑的具体行为:

  • 关联目标表silver和临时视图updates(当前微批数据)
  • 以s.mrn = u.mrn为匹配条件判断记录是否存在
  • 当匹配且任意业务字段(如出生日期dob、姓名first_name等)或更新时间updated不一致时,更新Silver表的对应记录
  • 当不匹配时,直接插入当前微批的新记录

完整代码

def upsert_to_delta(microBatchDF, batchId):
    microBatchDF.createOrReplaceTempView("updates")
    microBatchDF._jdf.sparkSession().sql("""
        MERGE INTO silver s
        USING updates u
        ON s.mrn = u.mrn
        WHEN MATCHED AND s.dob <> u.dob OR
                         s.sex <> u.sex OR
                         s.gender <> u.gender OR
                         s.first_name <> u.first_name OR
                         s.last_name <> u.last_name OR
                         s.street_address <> u.street_address OR
                         s.zip <> u.zip OR
                         s.city <> u.city OR
                         s.state <> u.state OR
                         s.updated <> u.updated
            THEN UPDATE SET *
        WHEN NOT MATCHED
            THEN INSERT *
    """)

query = (spark.readStream
              .table("bronze")
              .writeStream
              .foreachBatch(upsert_to_delta)
              .trigger(availableNow=True)
            #   .trigger(processingTime='5 seconds')
              .start())

DA.block_until_stream_is_ready(query)

内容的提问来源于stack exchange,提问作者Sajid Sarkar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 17:42:45