【发布时间】: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