【发布时间】:2017-05-28 00:11:14
【问题描述】:
我知道 PySpark DataFrames 是不可变的,所以我想创建一个新列,该列由应用于 PySpark DataFrame 现有列的转换产生。我的数据太大,无法使用 collect()。
所讨论的列是唯一整数列表的列表(给定列表中没有重复整数),例如:
[1]
[1,2]
[1,2,3]
[2,3]
上面是一个玩具示例,因为我的实际 DataFrame 具有最大长度为 52 个唯一整数的列表。我想生成一个列,它遍历整数列表并为每个循环删除一个元素。要删除的元素将是所有列表中唯一元素集合中的一个,在本例中为 [1,2,3]。
所以对于第一次迭代:
删除元素 1,结果如下:
[]
[2]
[2,3]
[2,3]
第二次迭代:
删除元素 2,结果如下:
[1]
[1]
[1,3]
[3]
等等。并重复上面的元素 3。
对于每次迭代,我想将结果附加到原始 PySpark DataFrame 以运行一些查询,使用这个“过滤”列作为原始 DataFrame 的行过滤器。
我的问题是,如何将 PySpark DataFrame 的列转换为列表?我的数据集很大,所以df.select('columnofintlists').collect() 会导致内存问题(例如:Kryo serialization failed: Buffer overflow. Available: 0, required: 1448662. To avoid this, increase spark.kryoserializer.buffer.max value.)。
【问题讨论】:
标签: pyspark