【问题标题】:Using pandas functions with Pyspark在 Pyspark 中使用 pandas 函数
【发布时间】:2021-09-12 14:19:54
【问题描述】:

我正在尝试使用 Pyspark 重写我的 Python 脚本 (Pandas),但我找不到一种方法来应用我的 Pandas 函数以提高 Pyspark 函数的效率:

我的功能如下:

def decompose_id(id_flight):
    
    my_id=id_flight.split("_")
    Esn=my_id[0]
    Year=my_id[3][0:4]
    Month=my_id[3][4:6]

return Esn, Year, Month

def reverse_string(string):
  stringlength=len(string) # calculate length of the list
  slicedString=string[stringlength::-1] # slicing 
  return slicedString

我想将第一个函数应用于数据框的一列(在 Pandas 中,我得到一行三个元素) 第二个函数用于验证 DataFrame 列的条件时使用

有没有使用 Pyspark 数据框应用它们的方法?

【问题讨论】:

  • Pandas 和 Spark 的工作方式不同。请用示例输入和输出解释您想要做什么。忘记你的 pandas 函数,解释预期的行为。
  • 顺便说一句,reverse_string 应该只是def reverse_string(string):return string[::-1]。而string 是内置库的名称,最好使用另一个词,例如in_string
  • 感谢您的评论!

标签: python pandas pyspark bigdata user-defined-functions


【解决方案1】:

您可以将这些函数作为 UDF 应用到 Spark 列,但效率不高。

以下是您执行任务所需的功能:

  • reverse :用它来替换你的函数reverse_string
  • split :用于替换 my_id=id_flight.split("_")
  • getItem :使用它来获取拆分列表中的项目my_id[3]
  • substr:替换python中的切片[0:4]

只需组合这些 spark 函数即可重新创建相同的行为。

【讨论】:

  • 效率不高是什么意思? (你指的是火花性能吗?)
  • @f.ivy 在 PySpark 中运行 UDF 是一个相当大的性能问题。当我们运行 UDF 时,Spark 需要对数据进行序列化,将其从 Spark 进程传输到 Python,对其进行反序列化,运行函数,对结果进行序列化,将其从 Python 进程移回 Scala,然后进行反序列化。
  • 我明白了,谢谢史蒂文的澄清
【解决方案2】:

如果您想利用 pandas 功能,一种方法是使用 - Pandas APIgroupBy

它为您提供了一种将每个 groupBy 集视为可以在其上实现功能的 pandas 数据框的方法。

然而,自从它使用 Spark 以来,架构实施是非常必要的,因为您还将浏览链接中提供的示例

一个实现示例可以找到here

对于琐碎的任务,比如反转字符串,选择内置 Spark 函数,否则使用 UDF

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-08
    • 2020-03-27
    • 1970-01-01
    • 1970-01-01
    • 2022-07-08
    相关资源
    最近更新 更多