请说明代码中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
相关产品推荐
相关产品推荐

