【问题标题】:spark structured streaming joining aggregate dataframe to dataframe火花结构化流将聚合数据帧连接到数据帧
【发布时间】:2018-11-08 07:18:55
【问题描述】:

我有一个流式数据框,它可能看起来像:

+--------------------+--------------------+
|               owner|              fruits|
+--------------------+--------------------+
|Brian                | apple|
Brian                | pear |
Brian                | date|
Brian                | avocado|
Bob                | avocado|
Bob                | apple|
........
+--------------------+--------------------+

我执行了 groupBy, agg collect_list 来清理。

val myFarmDF = farmDF.withWatermark("timeStamp", "1 seconds").groupBy("fruits").agg(collect_list(col("fruits")) as "fruitsA")

输出是每个所有者的单行和每个水果的数组。 我现在想将这个清理后的数组加入到原始流数据帧中,删除 fruits col 并且只有 fruitsA 列

val joinedDF = farmDF.join(myFarmDF, "owner").drop("fruits")

这似乎在我的脑海中起作用,但 spark 似乎不同意。

我得到一个

Failure when resolving conflicting references in Join:
'Join Inner
...
+- AnalysisBarrier
      +- Aggregate [name#17], [name#17, collect_list(fruits#61, 0, 0) AS fruitA#142]

当我把所有东西都变成一个静态数据框时,它工作得很好。这在流式上下文中是不可能的吗?

【问题讨论】:

    标签: scala apache-spark spark-structured-streaming


    【解决方案1】:

    您是否尝试过重命名列名?有类似问题https://issues.apache.org/jira/browse/SPARK-19860

    【讨论】:

    • 我居然想通了,忘记更新帖子了。
    • @Brian 这段代码对你有用吗?我将 Kafka 用于源和接收器,如果使用 outputMode("append"),我会得到 Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark;;,而如果我使用 outputMode("update"),我会得到 Inner join between two streaming DataFrames/Datasets is not supported in Update output mode, only in Append output mode;;。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-07-08
    • 2019-11-26
    • 2019-05-26
    • 2021-10-23
    • 2020-01-21
    • 2022-01-23
    • 1970-01-01
    相关资源
    最近更新 更多