【问题标题】:Find threshold in pyspark dataframe data在 pyspark 数据帧数据中查找阈值
【发布时间】:2022-01-08 10:01:56
【问题描述】:

总的来说,我对数据帧和 pyspark 完全陌生。

在 python 中,我想做的只是微不足道的 - 但是我似乎找不到使用 pyspark 不需要很长时间的方法。

我有一个大约 4000 行的 pyspark 数据框,其架构如下:

root
 |-- waveformData: struct (nullable = true)
 |    |-- elements: array (nullable = true)
 |    |    |-- element: double (containsNull = true)
 |    |-- dimensions: array (nullable = true)
 |    |    |-- element: integer (containsNull = true)

每个数组大约有 20000 个双精度数。

搜索此数组以查找最大值和阈值(最大值的 50% 的第一个实例)只需要很少的时间 - 但只有在数据采用“正常”格式时(numpy 数组)。

我使用的是基本款:

wav_df = temp_data.select("waveformData").toPandas()
wav = wav_df.to_numpy()[0][0].get("elements")

然后搜索最大值/阈值

但是 'toPandas' 步骤需要很长时间(比如单行需要 30 秒)

为什么?

我一直在尝试使用 .collect 等对 pyspark 数据帧进行操作以避免这种转换,但我尝试的一切都需要很长时间。

如果 pyspark 是针对大数据的,那我肯定做错了,这不可能是处理这么多数据的正常时间。

我错过了什么?

【问题讨论】:

  • 能否也为waveformData添加一些测试数据,这样我就可以使用该测试数据创建一个数据框,看看可以做什么
  • 怎么加最好,20000元素的数组贴不上去。
  • 我想我发现了这个问题,我一次加载一行并转换为熊猫,如果我一次转换所有 4000 行似乎要快得多,但后来我内存不足:(
  • 失败并显示以下错误消息:原因:org.apache.spark.SparkException:作业因阶段失败而中止:1286 个任务 (1025.2 MiB) 的序列化结果的总大小大于 spark。 driver.maxResultSize (1024.0 MiB)
  • 我尝试使用:spark.driver.maxResultSize g 但是得到'SparkSession'对象没有属性'驱动程序'

标签: pandas dataframe pyspark


【解决方案1】:

我创建了一些随机测试数据

from pyspark.sql import SparkSession
from pyspark.sql import Row
from pyspark.sql import functions as F

import numpy as np
np.random.seed(123)

spark = SparkSession.builder.getOrCreate()


n = 20000
for i in range(100):
    df = spark.createDataFrame([
        Row(waveFormData=Row(elements=[float(v) for v in np.random.randn(n)], dimensions=[n])) for i in range(40)
    ])
    df.write.parquet('waveFormData.parquet', mode='append')

当我加载数据并选择在 2 秒内运行的数组的最大值时:

df = spark.read.parquet('waveFormData.parquet')

df.select(F.array_max('waveFormData.elements')).toPandas()

【讨论】:

  • 好的,所以写成 parquet,然后读回并转换为 pandas?我会试试的。非常感谢。这也能解决我的内存问题吗?从我发现它正在将集群中的所有数据拉入驱动程序并使用太多内存。使用 Parquet 能解决这个问题吗?
  • 当您使用.toPandas() 方法时,您已从 Spark-Dataframe 切换到 Pandas-Dataframe。后者没有选择方法。试试up_data_par.select(F.array_max('waveformData.elements').alias('maxs')).select(F.max('maxs')).toPandas()
  • 确切的行为取决于您的数据源,但一般来说:是的,它很可能是随机排序的。你真的应该得到某种 ID 列,你可以用它来连接两个数据框。
  • 这种情况下,还要选择时间戳列,最后排序。不一定会保留开头的排序。或者用另一个 DF 做一个实际的.join(),这将是“正确”的做法。 :)
  • 我没有,但可能有办法。我对嵌套数据框不是很熟悉。我建议您展平您的表格(即将其转换为包含列 array_date、array_id、element_id、value 的平面表格),然后使用“普通”Spark SQL 函数(可能包括 groupby 和 Window 函数)进行分析
猜你喜欢
  • 1970-01-01
  • 2021-08-18
  • 2020-12-21
  • 2015-05-26
  • 1970-01-01
  • 1970-01-01
  • 2021-03-22
  • 2018-05-05
  • 1970-01-01
相关资源
最近更新 更多