【问题标题】:pysprak - microbatch streaming delta table as a source to perform merge against another delta table - foreachbatch is not getting invokedspark - 微批处理流式增量表作为对另一个增量表执行合并的源 - 未调用 foreachbatch
【发布时间】:2021-02-12 06:33:52
【问题描述】:

我已经创建了一个增量表,现在我正在尝试使用 foreachBatch() 将数据合并到该表中。我关注了这个example。我在谷歌云中的 dataproc image 1.5x 中运行此代码。

Spark 版本 2.4.7 Delta 版本 0.6.0

我的代码如下:

from delta.tables import *

spark = SparkSession.builder \
    .appName("streaming_merge") \
    .master("local[*]") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()
  
# Function to upsert `microBatchOutputDF` into Delta table using MERGE
def mergeToDelta(microBatchOutputDF, batchId): 
  (deltaTable.alias("accnt").merge(                                            
       microBatchOutputDF.alias("updates"), \
       "accnt.acct_nbr = updates.acct_nbr") \
     .whenMatchedDelete(condition = "updates.cdc_ind='D'") \
     .whenMatchedUpdateAll(condition = "updates.cdc_ind='U'") \
     .whenNotMatchedInsertAll(condition = "updates.cdc_ind!='D'") \
     .execute()
  )

deltaTable = DeltaTable.forPath(spark, "gs:<<path_for_the_target_delta_table>>")

# Define the source extract
SourceDF = (
  spark.readStream \
    .format("delta") \
    .load("gs://<<path_for_the_source_delta_location>>")

# Start the query to continuously upsert into target tables in update mode
SourceDF.writeStream \
  .format("delta") \
  .outputMode("update") \
  .foreachBatch(mergeToDelta) \
  .option("checkpointLocation","gs:<<path_for_the_checkpint_location>>") \
  .trigger(once=True) \
  .start() \

这段代码运行没有任何问题,但是没有数据写入增量表,我怀疑 foreachBatch 没有被调用。有谁知道我做错了什么?

【问题讨论】:

  • 如果你之前已经运行过你的代码,那么最后一个位置会存储在检查点中,直到上游发生变化,你不会得到新的变化......如果你想重新处理所有内容,尝试删除检查点。另外,将日志记录添加到 foreachbatch 函数
  • 当我在 spark shell 中单独执行相同的步骤时,我看到 foreachBatch 被调用。但是,当我使用 spark-shell .py 在终端中运行脚本时,它没有调用。只是不确定我是否做对了。我试过 1. gcloud dataproc 作业提交 pyspark --cluster=xxxxx --region=xxxx gs://>.py 2. 直接对 dataproc 集群执行 ssh 并使用“spark-submit gs:/ 触发” />.py"

标签: pyspark delta-lake


【解决方案1】:

添加 awaitTermination 后,流式传输开始工作并从源中获取最新数据并在 delta 目标表上执行合并。

【讨论】:

    猜你喜欢
    • 2020-12-23
    • 2023-01-10
    • 2019-09-21
    • 2017-12-18
    • 1970-01-01
    • 2019-01-14
    • 1970-01-01
    • 1970-01-01
    • 2023-01-17
    相关资源
    最近更新 更多