代码之家  ›  专栏  ›  技术社区  ›  Bryce Ramgovind

PySpark-按获取分组中每个列表的大小

  •  1
  • Bryce Ramgovind  · 技术社区  · 8 年前

    我有一个巨大的pyspark数据框架。我需要分组 Person 然后 collect 他们的 Budget 项目,以执行进一步的计算。 例如,

    a = [('Bob', 562,"Food", "12 May 2018"), ('Bob',880,"Food","01 June 2018"), ('Bob',380,'Household'," 16 June 2018"),  ('Sue',85,'Household'," 16 July 2018"), ('Sue',963,'Household'," 16 Sept 2018")]
    df = spark.createDataFrame(a, ["Person", "Amount","Budget", "Date"])
    

    分组依据:

    import pyspark.sql.functions as F
    df_grouped = df.groupby('person').agg(F.collect_list("Budget").alias("data"))
    

    架构:

    root
     |-- person: string (nullable = true)
     |-- data: array (nullable = true)
     |    |-- element: string (containsNull = true)
    

    然而,当我尝试对每个人应用UDF时,我会遇到一个内存错误。如何获取每个列表的大小(兆字节或千兆字节)( data )对于每个人?

    我已经做了以下工作,但我正在 nulls

    import sys
    size_list_udf = F.udf(lambda data: sys.getsizeof(data)/1000, DoubleType())
    df_grouped = df_grouped.withColumn("size",size_list_udf("data") )
    df_grouped.show()
    

    输出:

    +------+--------------------+----+
    |person|                data|size|
    +------+--------------------+----+
    |   Sue|[Household, House...|null|
    |   Bob|[Food, Food, Hous...|null|
    +------+--------------------+----+
    
    1 回复  |  直到 8 年前
        1
  •  1
  •   pault Tanjin    8 年前

    您的代码只有一个小问题。 sys.getsizeof() 以整数形式返回对象的大小(以字节为单位)。你要用这个除以整数值 1000 获取KB。在python 2中,这将返回一个整数。无论您如何定义 udf 返回 DoubleType() . 简单的解决方法是除以 1000.0 .

    import sys
    size_list_udf = f.udf(lambda data: sys.getsizeof(data)/1000.0, DoubleType())
    df_grouped = df_grouped.withColumn("size",size_list_udf("data") )
    df_grouped.show(truncate=False)
    #+------+-----------------------+-----+
    #|person|data                   |size |
    #+------+-----------------------+-----+
    #|Sue   |[Household, Household] |0.112|
    #|Bob   |[Food, Food, Household]|0.12 |
    #+------+-----------------------+-----+
    

    我发现 自定义项 正在返回 null ,罪魁祸首往往是类型不匹配。