【发布时间】:2020-09-12 09:40:57
【问题描述】:
问题归结为以下几点:我想在 pyspark 中使用现有的并行化输入集合和一个给定一个输入的函数生成一个 DataFrame,该函数可以生成相对大量的行。在下面的示例中,我想使用例如生成 10^12 行数据框1000 名执行者:
def generate_data(one_integer):
import numpy as np
from pyspark.sql import Row
M = 10000000 # number of values to generate per seed, e.g. 10M
np.random.seed(one_integer)
np_array = np.random.random_sample(M) # generates an array of M random values
row_type = Row("seed", "n", "x")
return [row_type(one_integer, i, float(np_array[i])) for i in range(M)]
N = 100000 # number of seeds to try, e.g. 100K
list_of_integers = [i for i in range(N)]
list_of_integers_rdd = spark.sparkContext.parallelize(list_of_integers)
row_rdd = list_of_integers_rdd.flatMap(list_of_integers_rdd)
from pyspark.sql.types import StructType, StructField, FloatType, IntegerType
my_schema = StructType([
StructField("seed", IntegerType()),
StructField("n", IntegerType()),
StructField("x", FloatType())])
df = spark.createDataFrame(row_rdd, schema=my_schema)
(我真的不想研究给定种子的随机数分布 - 这只是我能够想出的一个例子来说明大型数据帧不是从仓库加载而是由代码生成的情况)
上面的代码几乎完全符合我的要求。问题是它以一种非常低效的方式进行 - 代价是为每一行创建一个 Python Row 对象,然后将 Python Row 对象转换为内部 Spark 列表示。
有没有一种方法我可以通过让 spark 知道这些是一批值的列来转换已经以列表示形式的一批行(例如,一个或几个如上 np_array 的 numpy 数组)?
例如我可以编写代码来生成 python 集合 RDD,其中每个元素都是 pyarrow.RecordBatch 或 pandas.DataFrame,但我找不到将其中任何一个转换为 Spark DataFrame 的方法,而无需在过程。
至少有十几篇文章提供了如何使用 pyarrow + pandas 将本地(到驱动程序)pandas 数据帧有效地转换为 Spark 数据帧的示例,但这对我来说不是一个选择,因为我实际上需要数据在执行程序上以分布式方式生成,而不是在驱动程序上生成一个 pandas 数据帧并将其发送给执行程序。
UPD。 我找到了一种避免创建 Row 对象的方法——使用 Python 元组的 RDD。正如预期的那样,它仍然太慢了,但仍然比使用 Row 对象快一点。尽管如此,这并不是我真正想要的(这是一种将列数据从 python 传递到 Spark 的非常有效的方法)。
还测量了在机器上执行某些操作的时间(粗略的方法,测量时间有相当多的变化,但在我看来它仍然具有代表性):
有问题的数据集是 10M 行,3 列(一列是常数整数,另一列是从 0 到 10M-1 的整数范围,第三是使用 np.random.random_sample 生成的浮点值:
- 本地生成 pandas 数据帧(10M 行):~440-450ms
- 本地生成 spark.sql.Row 对象的 python 列表(10M 行):~12-15s
- 本地生成表示行(10M 行)的 python 元组列表:~3.4-3.5s
仅使用 1 个执行程序和 1 个初始种子值生成 Spark 数据帧:
- 使用
spark.createDataFrame(row_rdd, schema=my_schema): ~70-80s - 使用
spark.createDataFrame(tuple_rdd, schema=my_schema): ~40-45s - (非分布式创建)使用
spark.createDataFrame(pandas_df, schema=my_schema):~0.4-0.5s(没有 pandas df 生成本身,这需要大致相同的时间) -spark.sql.execution.arrow.enabled设置为 true。
本地到驱动程序 pandas 数据帧在约 1 秒内转换为 Spark 数据帧的 1000 万行的示例让我有理由相信执行程序中生成的数据帧应该可以实现。然而,我现在可以达到的最快速度是使用 Python 元组的 RDD 处理 1000 万行约 40 秒。
所以问题仍然存在 - 有没有办法在 pyspark 中以分布式方式有效地生成大型 Spark 数据帧?
【问题讨论】:
标签: apache-spark pyspark pyarrow apache-arrow