【问题标题】:Parallel Processing on Pyspark of around 1000 columnsPyspark 上大约 1000 列的并行处理
【发布时间】:2020-10-02 21:30:04
【问题描述】:

我有一个大约 1500 列的数据集,我正在尝试将所有列的零替换为 Null。如何在 Pyspark 中有效地做到这一点?

我曾尝试使用 spark UDF,但由于数据很宽,我无法并行处理多个列。

【问题讨论】:

    标签: apache-spark machine-learning pyspark


    【解决方案1】:

    我们不需要 UDF。您可以使用 spark 内置函数 df.na.replace 来实现它。你可以找到更多关于它的信息here。我写了一个实现相同目的的简单示例。

     from pyspark.sql import functions as F
    
     df = sc.parallelize([(1, 0, 5), (1,2, 0), (0,4, 5),  (1,7, 0), (0,0, 3),  
     (2,0, 5),  (2,3, 0)]).toDF(["a", "b", "c"])
    
            +---+---+---+
        |  a|  b|  c|
        +---+---+---+
        |  1|  0|  5|
        |  1|  2|  0|
        |  0|  4|  5|
        |  1|  7|  0|
        |  0|  0|  3|
        |  2|  0|  5|
        |  2|  3|  0|
        +---+---+---+
    
        df1=df.na.replace(0,None).show()
    
        +----+----+----+
        |   a|   b|   c|
        +----+----+----+
        |   1|null|   5|
        |   1|   2|null|
        |null|   4|   5|
        |   1|   7|null|
        |null|null|   3|
        |   2|null|   5|
        |   2|   3|null|
        +----+----+----+
    

    计算df中的不同值

        from pyspark.sql import functions as F
        df2=df1.agg(*(F.countDistinct(F.col(c)).alias(c) for c in df.columns))
    
        df2.show()
    
        +---+---+---+
        |  a|  b|  c|
        +---+---+---+
        |  2|  4|  2|
        +---+---+---+     
    

    要计算 99% 和 1%。

        df1.summary('99%', '1%').show()
    
        +-------+---+---+---+
        |summary|  a|  b|  c|
        +-------+---+---+---+
        |    99%|  2|  7|  5|
        |     1%|  1|  2|  3|
        +-------+---+---+---+
    

    【讨论】:

    • 谢谢。我可能提出了一个简单的问题。我还在考虑计算 df.summary() 提供之外的基本统计信息,例如计算不同的值、99%、1% 等。有没有一种有效的方法来做到这一点?
    猜你喜欢
    • 1970-01-01
    • 2011-06-11
    • 1970-01-01
    • 1970-01-01
    • 2021-11-14
    • 2013-11-19
    • 1970-01-01
    • 1970-01-01
    • 2013-02-16
    相关资源
    最近更新 更多