【问题标题】:Databricks, Question about "foreachBatch" to remove duplicate records when streaming data?Databricks,关于“foreachBatch”在流式传输数据时删除重复记录的问题?
【发布时间】: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


    【解决方案1】:

    Q1)要在流式传输时实现逻辑合并,您需要针对每个微批次执行此操作,因此可以在使用 foreachbatch API 时实现。

    Q2)您使用 dropDuplicates 和 Watermark 30 秒,如果您希望只能在 30 秒窗口(或您能够精确定义的任何窗口)中创建重复项,那么您的流将被重复数据删除。 (将会发生的是创建的流的状态)

    Q3)在实践中,您的 foreach 批次是(抱歉更多类似伪代码的 scala):

    .foreachBatch{ (microBatchDF: DataFrame, batch: Long) => 
            microBatchDF.createOrReplaceTempView(self.update_temp)
            microBatchDF._jdf.sparkSession().sql(self.sql_query)
          }
    

    希望这个对你有帮助

    【讨论】:

      猜你喜欢
      • 2019-12-25
      • 1970-01-01
      • 2021-05-18
      • 1970-01-01
      • 1970-01-01
      • 2021-02-23
      • 1970-01-01
      • 1970-01-01
      • 2012-08-31
      相关资源
      最近更新 更多