【问题标题】:How to zip two array columns in Spark SQL如何在 Spark SQL 中压缩两个数组列
【发布时间】:2019-06-14 09:56:48
【问题描述】:

我有一个 Pandas 数据框。我尝试先将包含字符串值的两列连接到一个列表中,然后使用 zip,我用“_”连接了列表的每个元素。我的数据集如下:

df['column_1']: 'abc, def, ghi'
df['column_2']: '1.0, 2.0, 3.0'

我想将这两列加入第三列,如下所示,用于我的数据框的每一行。

df['column_3']: [abc_1.0, def_2.0, ghi_3.0]

我已经使用下面的代码在 python 中成功完成了此操作,但是数据帧非常大,并且需要很长时间才能为整个数据帧运行它。为了提高效率,我想在 PySpark 中做同样的事情。我已成功读取 spark 数据框中的数据,但我很难确定如何使用 PySpark 等效函数复制 Pandas 函数。如何在 PySpark 中获得我想要的结果?

df['column_3'] = df['column_2']
for index, row in df.iterrows():
  while index < 3:
    if isinstance(row['column_1'], str):      
      row['column_1'] = list(row['column_1'].split(','))
      row['column_2'] = list(row['column_2'].split(','))
      row['column_3'] = ['_'.join(map(str, i)) for i in zip(list(row['column_1']), list(row['column_2']))]

我已使用以下代码将两列转换为 PySpark 中的数组

from pyspark.sql.types import ArrayType, IntegerType, StringType
from pyspark.sql.functions import col, split

crash.withColumn("column_1",
    split(col("column_1"), ",\s*").cast(ArrayType(StringType())).alias("column_1")
)
crash.withColumn("column_2",
    split(col("column_2"), ",\s*").cast(ArrayType(StringType())).alias("column_2")
)

现在我只需要使用“_”压缩两列中数组的每个元素。我该如何使用 zip 呢?任何帮助表示赞赏。

【问题讨论】:

  • 为什么df['column_1']df['column_2'] 是单个字符串而不是项目列表?它们最初是什么?
  • 数据是这样的,我正在数据框中读取数据
  • @Falconic 所以 abcdef 等在单行或不同行?同样是第 2 列单行?
  • @anky_91 这是 column_1 和 column_2 的一行数据框。每行在一列中有多个项目。这就是我拆分字符串然后转换为列表的原因。
  • 这能回答你的问题吗? Pyspark: Split multiple array columns into rows

标签: python pandas apache-spark pyspark apache-spark-sql


【解决方案1】:

Python 的 Spark SQL 等价物是 pyspark.sql.functions.arrays_zip:

pyspark.sql.functions.arrays_zip(*cols)

集合函数:返回结构的合并数组,其中第 N 个结构包含输入数组的所有第 N 个值。

所以如果你已经有两个数组:

from pyspark.sql.functions import split

df = (spark
    .createDataFrame([('abc, def, ghi', '1.0, 2.0, 3.0')])
    .toDF("column_1", "column_2")
    .withColumn("column_1", split("column_1", "\s*,\s*"))
    .withColumn("column_2", split("column_2", "\s*,\s*")))

你可以把它应用到结果上

from pyspark.sql.functions import arrays_zip

df_zipped = df.withColumn(
  "zipped", arrays_zip("column_1", "column_2")
)

df_zipped.select("zipped").show(truncate=False)
+------------------------------------+
|zipped                              |
+------------------------------------+
|[[abc, 1.0], [def, 2.0], [ghi, 3.0]]|
+------------------------------------+

现在合并结果你可以transformHow to use transform higher-order function?TypeError: Column is not iterable - How to iterate over ArrayType()?):

df_zipped_concat = df_zipped.withColumn(
    "zipped_concat",
     expr("transform(zipped, x -> concat_ws('_', x.column_1, x.column_2))")
) 

df_zipped_concat.select("zipped_concat").show(truncate=False)
+---------------------------+
|zipped_concat              |
+---------------------------+
|[abc_1.0, def_2.0, ghi_3.0]|
+---------------------------+

注意

Apache Spark 2.4 中引入了高阶函数 transformarrays_zip

【讨论】:

  • 感谢 user10465355。该解决方案对我有用,但请注意。它不能很好地处理列表中的空值。在将它们连接在一起之前,我手动从两列中删除了空值。其次,我必须在原始数据框中完成每一步。具有相同列名的多个数据框在 PySpark 中效果不佳。我必须调试它以查看我的代码中有什么问题。事实证明,我需要在不同的操作中使用相同的数据框。
【解决方案2】:

你也可以用UDF压缩分割数组列,

df = spark.createDataFrame([('abc,def,ghi','1.0,2.0,3.0')], ['col1','col2']) 
+-----------+-----------+
|col1       |col2       |
+-----------+-----------+
|abc,def,ghi|1.0,2.0,3.0|
+-----------+-----------+ ## Hope this is how your dataframe is

from pyspark.sql import functions as F
from pyspark.sql.types import *

def concat_udf(*args):
    return ['_'.join(x) for x in zip(*args)]

udf1 = F.udf(concat_udf,ArrayType(StringType()))
df = df.withColumn('col3',udf1(F.split(df.col1,','),F.split(df.col2,',')))
df.show(1,False)
+-----------+-----------+---------------------------+
|col1       |col2       |col3                       |
+-----------+-----------+---------------------------+
|abc,def,ghi|1.0,2.0,3.0|[abc_1.0, def_2.0, ghi_3.0]|
+-----------+-----------+---------------------------+

【讨论】:

  • 谢谢@suresh。这绝对是一个更清洁的解决方案。当我将它应用到我自己的数据框并运行收集函数时,出现以下错误 TypeError: zip argument #1 must support iteration 有什么想法吗?
  • 错误是因为,zip() 没有获得可迭代的输入。您能否发布您的示例输入数据框及其架构。
【解决方案3】:

对于 Spark 2.4+,这可以通过仅使用 zip_with 函数同时压缩连接来完成:

df.withColumn("column_3", expr("zip_with(column_1, column_2, (x, y) -> concat(x, '_', y))")) 

高阶函数使用 lambda 函数 (x, y) -&gt; concat(x, '_', y) 逐元素合并 2 个数组。

【讨论】:

    【解决方案4】:

    假设您已经得到问题中提到的两列数组,

    PySpark 现在为 Spark 3.1+ 提供了direct function callzip_with(),因此可以这样做:

    import pyspark.sql.functions as F
    
    df = df.withColumn(
        "column_3", 
        F.zip_with(
            "column_1", "column_2", 
            lambda x,y: F.concat_ws("_", x, y)
        )
    )
    

    【讨论】:

      猜你喜欢
      • 2014-03-27
      • 1970-01-01
      • 2015-12-29
      • 1970-01-01
      • 2023-03-25
      • 2019-06-16
      • 2011-11-20
      相关资源
      最近更新 更多