【发布时间】:2019-12-24 10:51:02
【问题描述】:
我有一个监视器目录包含多个.csv 文件。我需要在每个即将到来的.csv 文件中计算number of entries。我想在 pyspark 流式传输上下文中执行此操作。
这就是我所做的,
my_DStream = ssc.textFileStream(monitor_Dir)
test = my_DStream.flatMap(process_file) # process_file function simply process my file. e.g line.split(";")
print(len(test.collect()))
这并没有给我想要的结果。例如 file1.csv 包含 10 条目,file2.csv 包含 18 条目等。所以我需要查看输出
10
18
..
..
etc
如果我有一个单独的静态文件并使用 rdd 操作,我可以执行相同的任务。
【问题讨论】:
标签: python-3.x pyspark bigdata spark-streaming rdd