【发布时间】:2016-10-13 15:58:13
【问题描述】:
我们有时间序列数据(自 1970 年以来的时间戳和整数数据值):
# load data and cache it
df_cache = readInData() # read data from several files (paritioned by hour)
df_cache.persist(pyspark.StorageLevel.MEMORY_AND_DISK)
df_cache.agg({"data": "max"}).collect()
# now data is cached
df_cache.show()
+--------------------+---------+
| time| data|
+--------------------+---------+
|1.448409599861109E15|1551.7468|
|1.448409599871109E15|1551.7463|
|1.448409599881109E15|1551.7468|
现在我们想使用外部 python 库在 10 分钟时间窗口之上计算一些重要的事情。为此,我们需要将每个时间帧的数据加载到内存中,应用外部函数并存储结果。因此,用户定义的聚合函数 (UDAF) 是不可能的。
现在的问题是,当我们将 GroupBy 应用于 RDD 时,它非常慢。
df_cache.rdd.groupBy(lambda x: int(x.time / 600e6) ). \ # create 10 minute groups
map(lambda x: 1). \ # do some calculations, e.g. external library
collect() # get results
此操作需要 14 分钟左右在两个具有 6GB Ram 的节点上进行 120Mio 样本(100Hz 数据)。 groupBy 阶段的 Spark 详细信息:
Total Time Across All Tasks: 1.2 h
Locality Level Summary: Process local: 8
Input Size / Records: 1835.0 MB / 12097
Shuffle Write: 1677.6 MB / 379
Shuffle Spill (Memory): 79.4 GB
Shuffle Spill (Disk): 1930.6 MB
如果我使用一个简单的 python 脚本并让它遍历输入文件,完成所需的时间会更少。
如何在 spark 中优化这项工作?
【问题讨论】:
-
我不确定这是否可行,但您可能希望将数据与 int(x.time, 600e^) 映射为元组的第三个成员,然后将 reduceByKey 映射到该成员.我不确定你可以通过这个获得多少性能,但它应该会产生更小的洗牌。让我知道你的情况。
-
类似的东西:
df_cache.rdd.map(lambda x: (int(x.time/600e^), x.time, x.data) ).reduceByKey(lambda x: 1).collect()(错过了编辑时间跨度......!)但要小心收集,这可能会导致 OOM 异常。 -
我以前试过这个。而且它比直接group by要慢。 `asdf
-
我几乎想不出比 reduceByKey 更快的方法了。您可能想查看 groupBy 生成的组的分布是否有任何异常,并查看 reduceByKey 作业的详细信息,以查明确切的问题操作。你试过用 Spark SQL 的 DataFrame 重写你的工作吗?
-
这只需 9.47 秒即可完成。似乎是对的,它更快:
df_cache.rdd.map(lambda x: (int(x.time / 600e6), (x.time, x.data)) ).reduceByKey(lambda x,y: 1).collect()
标签: time apache-spark pyspark series