【发布时间】:2017-10-12 07:19:14
【问题描述】:
我在名为part-0001、part-0002 等的 Linux 机器上的单个目录中有大约 200 个文件。每个都有大约一百万行具有相同的列(称它们为“a”、“b”等)。让 'a','b' 对成为每一行的键(有许多重复项)。
同时,我建立了一个 Spark 2.2.0 集群,有一个 master 和两个 slave,总共有 42 个可用内核。地址是spark://XXX.YYY.com:7077。
然后我使用 PySpark 连接到集群并计算每个唯一对的 200 个文件的计数,如下所示。
from pyspark import SparkContext
from pyspark.sql import SQLContext
import pandas as pd
sc = SparkContext("spark://XXX.YYY.com:7077")
sqlContext = SQLContext(sc)
data_path = "/location/to/my/data/part-*"
sparkdf = sqlContext.read.csv(path=data_path, header=True)
dfgrouped = sparkdf.groupBy(['a','b'])
counts_by_group = dfgrouped.count()
我可以看到 Spark 正在处理一系列消息,并且它确实返回了看起来合理的结果。
问题:在执行此计算时,top 没有显示任何证据表明从属内核正在执行任何操作。似乎没有任何并行化。每个从属设备在作业之前都有一个相关的 Java 进程(加上来自其他用户的进程和后台系统进程)。所以看起来主人正在做所有的工作。鉴于有 200 个奇怪的文件,我曾预计在每台从属计算机上运行 21 个进程,直到事情结束(这 是当我显式调用 parallelize 时看到的,如下 count = sc.parallelize(c=range(1, niters + 1), numSlices=ncores).map(f).reduce(add)一个单独的实现)。
问题:如何确保 Spark 实际并行化计数?我希望每个核心抓取一个或多个文件,对它在文件中看到的对执行计数,然后将单个结果缩减为单个 DataFrame。我不应该在顶部看到这个吗?我需要显式调用并行化吗?
(FWIW,我见过使用分区的示例,但我的理解是,这是用于在 single 文件的块上分配处理。我的情况是我有很多文件。)
提前致谢。
【问题讨论】:
标签: python apache-spark pyspark