【问题标题】:Spark streaming printing top-k results of a dstreamSpark 流式打印 dstream 的前 k 个结果
【发布时间】:2018-05-17 14:12:04
【问题描述】:

我正在研究 hadoop 平台,我正在试验的东西是 Spark-Streaming API。我正在尝试读取文件流以计算每 x 秒后的字数(历史的累积总和)。现在我想将前 k 个单词打印到文件中。这是我想要做的:

# sort the dstream for current batch
sorted_counts = counts.transform(lambda rdd: rdd.sortBy(lambda x: x[1], ascending=False))

# get the top K values of each rdd from the transformed dstream
topK = sorted_counts.transform(lambda rdd: rdd.take(k))

我可以使用以下命令将输出打印到控制台/日志文件:

sorted_counts.pprint(k)

但问题是当我尝试使用以下方式将其打印到文件时:

topK.saveAsTextFiles(out_path)

或者即使我尝试将 topK 打印到控制台:

topK.pprint()

我收到以下错误,

AttributeError: 'list' 对象没有属性 '_jrdd'

我认为这是因为 rdd.take(k) 返回实际列表而不是 rdd。我该如何解决它?此外,我想为每个新计算的字数生成不同的文件...即每 x 秒生成一个新的输出文件(使用 saveAsTextFiles() 保证。如果有帮助,我正在使用 python 进行编程。谢谢!

【问题讨论】:

  • sorted_counts.transform(lambda rdd: sc.parallelize(rdd.take(k)))

标签: python apache-spark pyspark


【解决方案1】:

似乎没有 API 可以让您这样做。但是您可以解决方法:

rdd.zipWithIndex().filter(<filter with big index>).map(<remove index here>)

另一种解决方案是(无排序):

sc.parallelize(rdd.top(...))

这样你就不需要对所有的 RDD 进行排序,只需要取出最大的元素然后从它们创建 RDD。

【讨论】:

  • 我已经编辑了我的答案 - 我认为添加了一个更好的解决方案。
  • 感谢您的回答。有效。第一个就是。将尝试第二个,看看会发生什么。
猜你喜欢
  • 2016-09-29
  • 2015-03-18
  • 2016-11-11
  • 1970-01-01
  • 1970-01-01
  • 2014-09-21
  • 2015-11-10
  • 2018-08-18
  • 2016-04-26
相关资源
最近更新 更多