【发布时间】:2021-09-12 17:26:59
【问题描述】:
我已经开始处理分布在几个类中的具有以下结构的代码:
var data1 = sparkSession.read().parquet("s3a://....");
log.info(..., data1.count())
var data2 = data1.filter(....)
log.info(..., data2.count())
var data3 = data2.filter(....)
log.info(..., data3.count())
var data4 = firstSubset(data3).union(secondSubset(data3))
log.info(...., data4.count())
data4.foreachPartition(... code to write to SQS ....)
firstSubset() 和 secondSubset() 执行更多过滤/计数,其中一个连接到从 S3 加载的另一个 parquet 文件,以通过内部连接进行过滤。该代码在 Amazon EMR 上使用 Spark 2.4 的 Java 实现,数据由 Spark Dataset<Row> 对象组成。目的是跟踪每个过滤器对原始数据的影响。
对于在本地运行的极少数行,此代码花费的时间过长(约 10 分钟)。我是 Spark 的新手,但我的理解是,每次我们有像 .count() 这样的操作时,Spark 都会创建一个执行计划并运行它 - 导致重复的数据加载、工作和节点之间的混洗。
我有这些问题:
- 我对 Spark 必须如何重新完成工作以完成每次计数的理解是否正确?
- 这会导致数据被多次下载吗?
- Spark 中是否有一种更有效的机制可以像这样从同一数据集中查找多个输出(即部分计数)?对我来说,这似乎很直观,如果不是为了与另一个 parquet 文件连接,这可以通过一次数据传递来完成。
- Spark 中是否有其他方法可以获取中间结果的计数?
【问题讨论】:
标签: java apache-spark amazon-emr apache-spark-dataset