【发布时间】:2019-10-27 12:35:57
【问题描述】:
我刚刚使用 GCP(dataproc) 和 pyspark 了解了这个美妙的大数据和云技术世界。我有 ~5 GB 大小的 JSON 文件(压缩,gz 文件),其中包含 ~500 万 条记录,我只需要读取每一行并处理满足特定条件的那些行。我有我的工作代码,我发出了一个带有 --num-partitions=5 的 spark-submit,但仍然只有一名工人用于执行该操作。
这是我正在使用的 spark-submit 命令:
spark-submit --num-executors 5 --py-files /home/user/code/dist/package-0.1-py3.6.egg job.py
job.py:
path = "gs://dataproc-bucket/json-files/data_5M.json.gz"
mi = spark.read.json(path)
inf_rel = mi.select(mi.client_id,
mi.user_id,
mi.first_date,
F.hour(mi.first_date).alias('hour'),
mi.notes).rdd.map(foo).filter(lambda x: x)
inf_relevance = inf_rel.map(lambda l: Row(**dict(l))).toDF()
save_path = "gs://dataproc-bucket/json-files/output_5M.json"
inf_relevance.write.mode('append').json(save_path)
print("END!!")
Dataproc 配置: (我现在使用免费帐户,一旦我得到工作解决方案将添加更多核心和执行器)
(Debian 9、Hadoop 2.9、Spark 2.4) 主节点:2 个 vCPU,7.50 GB 内存 主磁盘大小:32 GB 5 个工作节点:1 个 vCPU,3.75 GB 内存 主磁盘类型:32 GB
在 spark-submit 之后,我可以在 web UI 中看到添加了 5 个执行程序,但只有 1 个执行程序保持活动状态并执行所有任务,其余 4 个被释放。
我进行了研究,大部分问题都是关于通过 JDBC 访问数据的。
请提出我在这里缺少的内容。
附:最终我会读取 64 个 5 GB 的 json 文件,因此可能会使用 8 个核心 * 100 个工作人员。
【问题讨论】:
-
你的数据框有多少个分区?检查
df.rdd.getNumPartitions()...我没有使用这个google dataproc东西但是当从jdbc读取数据时它默认为1个分区,因为它是单线程的。我敢打赌你的dataframe是 1 个分区 = 1 个任务 = 1 个正在使用的核心,在这种情况下,只在一台机器上使用,因为没有任何东西被并行化。 -
谢谢,我现在将 CPU 数量增加到 2 X 3 个工作程序,并将其重新分区为 10,但仍然没有改进。我忘记提到的另一件重要的事情是这个 json 文件是压缩文件 (.gz) 格式。不确定这是否会导致任何问题。
-
是的,这是个问题... gzip 格式不可拆分... 可能与您的问题的解决方案在这里重复 => stackoverflow.com/questions/40492967/…
标签: json apache-spark pyspark google-cloud-storage google-cloud-dataproc