【问题标题】:How to drop duplicates while streaming in spark如何在 Spark 中流式传输时删除重复项
【发布时间】: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


【解决方案1】:

我发现您的代码有 2 个问题:

  1. 当您在函数 upsertToDelta 中调用 Merge 语句时。 Spark 需要一个可以合并“更新”tempView 的目标表。

在代码中:

合并到事件DF t 使用更新 ON s.deviceId = t.deviceId 不匹配时插入 *

eventsDF 应该是目标表名。

  1. 上述代码本身会将 tempView 与目标表连接起来,并在不匹配时插入到目标表中。

因此,不需要 start() 中的 writePath。 .start(writePath)

请注意:如果需要,您还可以在合并代码中添加更新选项。

【讨论】:

    猜你喜欢
    • 2018-07-22
    • 1970-01-01
    • 1970-01-01
    • 2014-07-11
    • 2018-05-27
    • 1970-01-01
    • 1970-01-01
    • 2021-01-05
    • 2020-07-17
    相关资源
    最近更新 更多