【发布时间】:2021-05-20 03:41:12
【问题描述】:
我有一个流式传输作业,将数据流式传输到 databricks spark 中的 delta 湖中,我试图在流式传输时删除重复项,因此我的 delta 数据没有重复项。到目前为止,这是我所拥有的:
inputPath = "my_input_path"
schema = StructType("some_schema")
eventsDF = (
spark
.readStream
.schema(schema)
.option("header", "true")
.option("maxFilesPerTrigger", 1)
.csv(inputPath)
)
def upsertToDelta(eventsDF, batchId):
eventsDF.createOrReplaceTempView("updates")
eventsDF._jdf.sparkSession().sql("""
MERGE INTO eventsDF t
USING updates s
ON s.deviceId = t.deviceId
WHEN NOT MATCHED THEN INSERT *
""")
writePath = "my_write_path"
checkpointPath = writePath + "/_checkpoint"
deltaStreamingQuery = (eventsDF
.writeStream
.format("delta")
.foreachBatch(upsertToDelta)
.option("checkpointLocation", checkpointPath)
.outputMode("append")
.queryName("test")
.start(writePath)
)
我收到错误消息:py4j.protocol.Py4JJavaError: An error occurred while calling o398.sql.
: org.apache.spark.sql.AnalysisException: Table or view not found: eventsDF; line 2 pos 4
但我刚刚开始流式传输这些数据,还没有创建任何表。
【问题讨论】:
-
你有什么问题?您是在询问未找到表或如何删除重复项?
-
@JacekLaskowski 如何删除重复项
-
那你为什么不用
Dataset.dropDuplicates呢? -
@JacekLaskowski,我正在尝试在流式传输数据时删除重复项,我的表不断有新数据,只能执行一次 dropDuplicates。
-
您将删除重复的内容,直到出现水印为止。
Dataset是流式查询(不是结构化的一次性批量查询)。你试过了吗?
标签: apache-spark databricks spark-structured-streaming delta-lake