【问题标题】:Combine PySpark DataFrame ArrayType fields into single ArrayType field将 PySpark DataFrame ArrayType 字段组合成单个 ArrayType 字段
【发布时间】:2016-09-14 00:33:21
【问题描述】:

我有一个带有 2 个 ArrayType 字段的 PySpark DataFrame:

>>>df
DataFrame[id: string, tokens: array<string>, bigrams: array<string>]
>>>df.take(1)
[Row(id='ID1', tokens=['one', 'two', 'two'], bigrams=['one two', 'two two'])]

我想将它们组合成一个 ArrayType 字段:

>>>df2
DataFrame[id: string, tokens_bigrams: array<string>]
>>>df2.take(1)
[Row(id='ID1', tokens_bigrams=['one', 'two', 'two', 'one two', 'two two'])]

适用于字符串的语法在这里似乎不起作用:

df2 = df.withColumn('tokens_bigrams', df.tokens + df.bigrams)

谢谢!

【问题讨论】:

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


    【解决方案1】:

    火花 >= 2.4

    你可以使用concat函数(SPARK-23736):

    from pyspark.sql.functions import col, concat 
    
    df.select(concat(col("tokens"), col("tokens_bigrams"))).show(truncate=False)
    
    # +---------------------------------+                                             
    # |concat(tokens, tokens_bigrams)   |
    # +---------------------------------+
    # |[one, two, two, one two, two two]|
    # |null                             |
    # +---------------------------------+
    

    要在其中一个值为 NULL 时保留数据,您可以使用 coalescearray

    from pyspark.sql.functions import array, coalesce      
    
    df.select(concat(
        coalesce(col("tokens"), array()),
        coalesce(col("tokens_bigrams"), array())
    )).show(truncate = False)
    
    # +--------------------------------------------------------------------+
    # |concat(coalesce(tokens, array()), coalesce(tokens_bigrams, array()))|
    # +--------------------------------------------------------------------+
    # |[one, two, two, one two, two two]                                   |
    # |[three]                                                             |
    # +--------------------------------------------------------------------+
    

    火花

    不幸的是,在一般情况下连接 array 列需要 UDF,例如:

    from itertools import chain
    from pyspark.sql.functions import col, udf
    from pyspark.sql.types import *
    
    
    def concat(type):
        def concat_(*args):
            return list(chain.from_iterable((arg if arg else [] for arg in args)))
        return udf(concat_, ArrayType(type))
    

    可以用作:

    df = spark.createDataFrame(
        [(["one", "two", "two"], ["one two", "two two"]), (["three"], None)], 
        ("tokens", "tokens_bigrams")
    )
    
    concat_string_arrays = concat(StringType())
    df.select(concat_string_arrays("tokens", "tokens_bigrams")).show(truncate=False)
    
    # +---------------------------------+
    # |concat_(tokens, tokens_bigrams)  |
    # +---------------------------------+
    # |[one, two, two, one two, two two]|
    # |[three]                          |
    # +---------------------------------+
    

    【讨论】:

    • 合并两个数组后如何删除重复项?
    • @j' df.withColumn('concat_no_duplicates', array_distinct(col('concat_(tokens, tokens_bigrams)')))
    【解决方案2】:

    在 Spark 2.4.0(Databricks 平台上的 2.3)中,您可以使用 concat 函数在 DataFrame API 中本地执行此操作。在您的示例中,您可以这样做:

    from pyspark.sql.functions import col, concat
    
    df.withColumn('tokens_bigrams', concat(col('tokens'), col('bigrams')))
    

    Here 是相关的jira。

    【讨论】:

      【解决方案3】:

      我使用的是 Spark

      from pyspark.sql import functions as F
      
      df.select("*",F.array(F.concat_ws(',', col('tokens'), col('bigrams))).\
                                  alias('concat_cols'))
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-12-25
        • 2020-01-08
        • 1970-01-01
        • 2018-10-21
        • 2020-10-02
        • 1970-01-01
        相关资源
        最近更新 更多