【问题标题】:Slow filtering of pyspark dataframespyspark 数据帧的慢速过滤
【发布时间】:2019-05-13 06:12:43
【问题描述】:

我对过滤 pandas 和 pyspark 数据帧时的时差有疑问:

import time
import numpy as np
import pandas as pd
from random import shuffle

from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()

df = pd.DataFrame(np.random.randint(1000000, size=400000).reshape(-1, 2))
list_filter = list(range(10000))
shuffle(list_filter)

# pandas is fast 
t0 = time.time()
df_filtered = df[df[0].isin(list_filter)]
print(time.time() - t0)
# 0.0072

df_spark = spark.createDataFrame(df)

# pyspark is slow
t0 = time.time()
df_spark_filtered = df_spark[df_spark[0].isin(list_filter)]
print(time.time() - t0)
# 3.1232

如果我将list_filter 的长度增加到 10000,那么执行时间是 0.01353 和 17.6768 秒。 isin seems 的 Pandas 实现具有计算效率。你能解释一下为什么 pyspark 数据帧的过滤如此缓慢,我怎样才能快速执行这种过滤?

【问题讨论】:

  • 您是否尝试增加分区大小并使用 spark-submit 运行?
  • @coldspeed 我是 spark 新手,能否详细说明一下?

标签: python pandas pyspark pyspark-sql


【解决方案1】:

您需要使用 join 代替带有 isin 子句的过滤器来加速 pyspark 中的过滤器操作:

import time
import numpy as np
import pandas as pd
from random import shuffle
import pyspark.sql.functions as F

from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()

df = pd.DataFrame(np.random.randint(1000000, size=400000).reshape(-1, 2))

df_spark = spark.createDataFrame(df)

list_filter = list(range(10000))
list_filter_df = spark.createDataFrame([[x] for x in list_filter], df_spark.columns[:1])
shuffle(list_filter)

# pandas is fast because everything in memory
t0 = time.time()
df_filtered = df[df[0].isin(list_filter)]
print(time.time() - t0)
# 0.0227580165863
# 0.0127580165863

# pyspark is slow because there is memory overhead, but broadcast make is mast compared to isin with lists
t0 = time.time()
df_spark_filtered = df_spark.join(F.broadcast(list_filter_df), df_spark.columns[:1])
print(time.time() - t0)
# 0.0571971035004
# 0.0471971035004

【讨论】:

  • 感谢您的回复!我认为,为了公平比较,我们需要添加df_spark_filtered.collect(),它更像是 2 秒而不是 50 毫秒。
【解决方案2】:

Spark 旨在用于处理大量数据。如果数据适合熊猫数据框,熊猫总是会更快。问题是,对于海量数据,pandas 会失败,而 spark 会完成这项工作(例如,比 MapReduce 更快)。

在这些情况下,Spark 通常较慢,因为它需要开发要执行的操作的 DAG,例如执行计划,并尝试对其进行优化。

所以,你应该只在数据很大的时候才考虑使用spark,否则使用pandas会更快。

您可以查看this article 并查看 pandas 和 spark 速度之间的比较,并且 pandas 总是更快,直到数据大到失败。

【讨论】:

  • 感谢您的回复。预计会产生开销,但开销太大了 + pyspark isinlist_filter 的扩展性非常差。
  • 尝试将spark.sql.shuffle.partitionsspark.default.parallelism 设置为1,以匹配pandas。
猜你喜欢
  • 2017-03-16
  • 2021-04-10
  • 2017-10-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-24
相关资源
最近更新 更多