【问题标题】:How can I combine a DStream pair of key and value using the same key?如何使用相同的键组合一对 DStream 键和值?
【发布时间】:2016-06-09 10:16:37
【问题描述】:

我想使用 spark 将第一个 DStream 更改为第二个。但我不知道该怎么做?我已经尝试过 groupByKey(),它不起作用和 aggregateByKey(),它只使用 RDD 而不是 DStream。

这是当前结果:

DStream [(1,value1),(2,value2),(3,value3),(1,value4),(1,value5),(2,value6)]

这是我想要的结果:

DStream(1,(value1,value4,value5)) ,(2,(value2,value5)) ,(3,(value3))

感谢您的回复。

【问题讨论】:

  • groupByKey 是什么意思没用
  • 将相同的键与 groupByKey 结合使用时不会给我相同的键和值对。这意味着它没有给我这个结果:DStream(1,(value1,value4,value5)),(2,(value2,value5)),(3,(value3))。我不知道该怎么做,也许我使用 groupByKey 错误?

标签: apache-spark spark-streaming datastax-enterprise


【解决方案1】:

groupByKey 正是这样做的。它将DStream[K, V] 转换为DStream[(K, Seq[V])]。我怀疑您对输出的期望可能是错误的。因为DStream 只是RDDs 组的无限序列,单独应用于每个RDD。所以如果第一批包含:

(1,value1),(2,value2),(3,value3),(1,value4)

第二个

(1,value5),(2,value6)

你会得到

(1, [value1, value4]), (2, [value2]), (3, value3)

(1,[value5]),(2,[value6])

分别。

虽然 DStreams 支持有状态操作 (updateStateByKey),但您不太可能希望将其用于不断增长的集合。

【讨论】:

  • 感谢您的解释。使用 groupByKey 不允许我组合相同的密钥对,因为它是一个流并使用多个 RDD/不断增长的集合。在我的情况下,您提出什么解决方案来实现上述结果?谢谢
  • 我不清楚你想要达到什么目的。我的意思是,流中不断增长的值迟早会破坏内存。如果您想保留所有内容(并且可能在需要时阅读)。如果您要查看更宽的间隔,请尝试窗口操作。
  • 我不确定 Spark Streaming 可以处理多少数据量以及何时用数据库分析替代 Stream。这是我更详细的问题:stackoverflow.com/questions/35691172/…。感谢您的回复和帮助!
  • 它认为除了单个数据结构的内存限制之外没有硬限制。可能在此之前很久,您就会遇到其他一些性能问题,例如密集溢出到磁盘或 GC。附带说明 - 如果这回答了问题的这一部分,请不要忘记投票/接受。
猜你喜欢
  • 2022-12-05
  • 1970-01-01
  • 2016-12-29
  • 1970-01-01
  • 2021-06-18
  • 2020-02-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多