【发布时间】: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