【问题标题】:PySpark and method chainingPySpark 和方法链
【发布时间】:2020-04-28 09:14:57
【问题描述】:

我是 PySpark 的新手,想知道您将如何在那里使用方法链。在 pandas 中,我会使用 assign 和 lambda,例如

df = pd.DataFrame({'number':[1,2,3],'date':['31-dec-19','02-jan-18','14-mar-20']})

df = (df.assign(number_plus_one = lambda x: x.number + 1)
        .assign(date = lambda x: pd.to_datetime(x.date))
        .loc[lambda x: x.number_plus_one.isin([2,3])]
        .drop(columns=['number','number_plus_one'])
      )

如何在 PySpark 中编写相同的代码而不将其转换为 pandas 数据帧?我猜你可以使用过滤器、withColumn 和 drop,但是你将如何使用方法链接呢?

【问题讨论】:

    标签: python pandas lambda pyspark


    【解决方案1】:

    您可以通过以类似方式链接调用在 Spark 中做同样的事情。这是一个例子:

    sc.parallelize([Row(number=1, date='31-dec-19'),
                         Row(number=1, date='31-dec-19'),
                         Row(number=1, date='31-dec-19')])\
    .toDF()\
    .withColumn('number_plus_one', f.col('number') + 1)\
    .filter(f.col('number_plus_one').cast(IntegerType()).isin(f.lit(2), f.lit(3)) )\
    .drop('number','number_plus_one')\
    .show()
    

    结果

    +---------+
    |     date|
    +---------+
    |31-dec-19|
    |02-jan-18|
    +---------+
    

    【讨论】:

    • 非常感谢。我不太了解您正在使用的f,我收到一个错误:An error was encountered: name 'f' is not defined。我该怎么办?
    • 对不起,这是别名:from pyspark.sql import functions as f
    • 太好了,谢谢。还有一个问题,pd.to_datetime(x.date) 你会怎么做?
    • 没有直接的方法来处理列值。您可以执行spark.createDataFrame([[1,'31-dec-19'],[2,'02-jan-18'],[3,'14-mar-20']], ['number', 'date']) 之类的操作。所以只需将你的列压缩成行。
    • 但是我还有日期字符串?我有一列包含日期但字符串不是日期格式。那我怎么把它转换成日期呢?
    猜你喜欢
    • 2011-11-06
    • 2016-05-22
    • 1970-01-01
    • 2021-02-06
    • 2021-08-18
    • 2011-07-28
    • 2021-02-07
    • 2016-02-10
    • 1970-01-01
    相关资源
    最近更新 更多