【问题标题】:Spark Time Series - Custom Group By 10 Minute Intervals: Improve PerformanceSpark 时间序列 - 按 10 分钟间隔自定义组:提高性能
【发布时间】: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


【解决方案1】:

groupBy 是您的瓶颈:它需要对所有分区的数据进行洗牌,这很耗时并且占用大量内存空间,从指标中可以看出。

这里的方法是使用reduceByKey 操作并将其链接如下: df_cache.rdd.map(lambda x: (int(x.time/600e6), (x.time, x.data) ).reduceByKey(lambda x,y: 1).collect()

这里的关键点是groupBy 需要在所有分区中对所有数据进行洗牌,而reduceByKey 将首先在每个分区上进行缩减,然后在所有分区上进行缩减 - 大大减少了全局洗牌的大小。请注意我如何将输入组织成一个键以利用 reduceByKey 操作。

正如我在 cmets 中提到的,您可能还想通过使用 Spark SQL 的 DataFrame 抽象来尝试您的程序,这可能会给您带来额外的提升,这要归功于它的优化器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-01-15
    • 2016-04-15
    • 1970-01-01
    • 1970-01-01
    • 2021-07-16
    • 2021-07-18
    • 2019-04-10
    • 2011-06-27
    相关资源
    最近更新 更多