【问题标题】:Remove element from PySpark DataFrame column从 PySpark DataFrame 列中删除元素
【发布时间】: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


    【解决方案1】:

    df.toLocalIterator() 将返回一个迭代器 for 循环

    【讨论】:

      【解决方案2】:

      这是 pyspark 文档中的一个示例:

      >>>from pyspark.sql.functions import array_remove
      >>>from pyspark.sql import SparkSession, SQLContext
      
      >>>sc = SparkContext.getOrCreate(SparkConf().setMaster("local[*]"))
      >>>spark = SparkSession(sc)
      
      >>>df = spark.createDataFrame([([1, 2, 3, 1, 1],), ([],)], ['data'])
      
      >>>df.select(array_remove(df.data, 1)).collect()
      [Row(array_remove(data, 1)=[2, 3]), Row(array_remove(data, 1)=[])]
      

      参考: https://spark.apache.org/docs/latest/api/python/pyspark.sql.html?highlight=w

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-01-14
        • 2018-12-31
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多