【发布时间】:2022-10-24 19:15:01
【问题描述】:
我正在练习使用 here 发布的 Databricks 示例笔记本:
在其中一个笔记本(ADE 3.1 - 流式重复数据删除)(URL) 中,有一个示例代码可以在流式传输数据时删除重复记录。
我对此有几个问题,希望能得到您的帮助。我复制以下代码的主要部分:
from pyspark.sql import functions as F
json_schema = "device_id LONG, time TIMESTAMP, heartrate DOUBLE"
deduped_df = (spark.readStream
.table("bronze")
.filter("topic = 'bpm'")
.select(F.from_json(F.col("value").cast("string"), json_schema).alias("v"))
.select("v.*")
.withWatermark("time", "30 seconds")
.dropDuplicates(["device_id", "time"]))
sql_query = """
MERGE INTO heart_rate_silver a
USING stream_updates b
ON a.device_id=b.device_id AND a.time=b.time
WHEN NOT MATCHED THEN INSERT *
"""
class Upsert:
def __init__(self, sql_query, update_temp="stream_updates"):
self.sql_query = sql_query
self.update_temp = update_temp
def upsert_to_delta(self, microBatchDF, batch):
microBatchDF.createOrReplaceTempView(self.update_temp)
microBatchDF._jdf.sparkSession().sql(self.sql_query)
streaming_merge = Upsert(sql_query)
query = (deduped_df.writeStream
.foreachBatch(streaming_merge.upsert_to_delta) # run query for each batch
.outputMode("update")
.option("checkpointLocation", f"{DA.paths.checkpoints}/recordings")
.trigger(availableNow=True)
.start())
query.awaitTermination()
Q1) 定义类Upsert 并使用方法foreachBatch 的原因是什么?
Q2) 如果我不使用foreachBatch 怎么办?
dropDuplicates(["device_id", "time"]) 方法在读取记录时删除重复项。确定没有重复记录还不够吗?
Q3)Upsert 类的方法upsert_to_delta 有两个输入参数(microBatchDF,batch)。但是,当我们在以下行中调用它时:
.foreachBatch(streaming_merge.upsert_to_delta)
,我们不传递它的论点。它如何获得 (microBatchDF, batch) 的值?
感谢您花时间阅读我的问题。
【问题讨论】:
标签: duplicates streaming databricks upsert