【问题标题】:Parse PySpark string column of key-list dictionaries based on separate array column of keys基于单独的键数组列解析键列表字典的 PySpark 字符串列
【发布时间】:2021-08-10 20:42:19
【问题描述】:

我目前正试图根据下面“Keys”列中的有序键数组提取下面“keyValue”列中的值。

>>> df.select('reference_nbr', 'keyValue', 'Keys').show()
+--------------+-------------------------------------------------------------+---------------------+
|    ref_number|                                                     keyValue|                 Keys|
+--------------+-------------------------------------------------------------+---------------------+
|          AZQ5|{key39=[TBAX3, TBAX6, TBAXN], key46=[TBARO, TBAZ4, TBABN],...|[key1, key2, key3,...|
|          NXY3|{key39=[TBAX8, TBAX2, TBAXZ], key46=[TBARD, TBAZK, TBAX9],...|[key1, key2, key3,...|
|          QSW6|{key39=[TBAX5, TBAX3, TBAX8], key46=[TBARB, TBAZN, TBAX4],...|[key1, key2, key3,...|
|          LJB7|{key39=[TBAX3, TBAXN, TBAXL], key46=[TBARM, TBAZ2, TBAX3],...|[key1, key2, key3,...|
|          MKH9|{key39=[TBAX4, TBAX9, TBAXV], key46=[TBARB, TBAZB, TBAX1],...|[key1, key2, key3,...|
|          UFG1|{key39=[TBAX3, TBAX6, TBAXQ], key46=[TBARL, TBAZB, TBAX0],...|[key1, key2, key3,...|
|          WDE4|{key39=[TBAX6, TBAX7, TBAX9], key46=[TBARX, TBAX6, TBAX8],...|[key1, key2, key3,...|
|          VRX8|{key39=[TBAX3, TBAX1, TBAX0], key46=[TBARQ, TBAX9, TBAX3],...|[key1, key2, key3,...|
|          CIZ2|{key39=[TBAX3, TBAXC, TBAX2], key46=[TBARA, TBAXQ, TBAX1],...|[key1, key2, key3,...|
|          BEO3|{key39=[TBAX9, TBAXQ, TBAX4], key46=[TBARP, TBAXV, TBAX2],...|[key1, key2, key3,...|
+--------------+-------------------------------------------------------------+---------------------+
only showing top 20 rows

如果我应用下面的 UDF 和 withColumn() 步骤,我可以轻松地根据特定键查询“keyValue”列,并将该键值的数组插入到新列中。

getKey4 = udf(lambda ar1: ar1.get('key4'))
df = df.withColumn("key4Values", getKey4(df["keyValue"]))

我正在尝试执行与上述相同的步骤,但按“键”列的顺序执行每个键。期望的输出:

+--------------+-------------------------------------------------------------+---------------------+------------------------------------------------+
|    ref_number|                                                     keyValue|                 Keys|                                          Values|
+--------------+-------------------------------------------------------------+---------------------+------------------------------------------------+
|          AZQ5|{key39=[TBAX3, TBAX6, TBAXN], key46=[TBARO, TBAZ4, TBABN],...|[key1, key2, key3,...|[[TBAX4, TBAXQ, TBAXD],[TBAR1, TBAZA, TBABW],...|
|          NXY3|{key39=[TBAX8, TBAX2, TBAXZ], key46=[TBARD, TBAZK, TBAX9],...|[key1, key2, key3,...|[[TBAX5, TBAXA, TBAXC],[TBAR2, TBAZS, TBABE],...|
|          QSW6|{key39=[TBAX5, TBAX3, TBAX8], key46=[TBARB, TBAZN, TBAX4],...|[key1, key2, key3,...|[[TBAX6, TBAXZ, TBAXF],[TBAR3, TBAZD, TBABR],...|
|          LJB7|{key39=[TBAX3, TBAXN, TBAXL], key46=[TBARM, TBAZ2, TBAX3],...|[key1, key2, key3,...|[[TBAX7, TBAXC, TBAXG],[TBAR4, TBAZF, TBABT],...|
|          MKH9|{key39=[TBAX4, TBAX9, TBAXV], key46=[TBARB, TBAZB, TBAX1],...|[key1, key2, key3,...|[[TBAX8, TBAXV, TBAXH],[TBAR5, TBAZG, TBABY],...|
|          UFG1|{key39=[TBAX3, TBAX6, TBAXQ], key46=[TBARL, TBAZB, TBAX0],...|[key1, key2, key3,...|[[TBAX9, TBAXB, TBAXJ],[TBAR6, TBAZH, TBABU],...|
|          WDE4|{key39=[TBAX6, TBAX7, TBAX9], key46=[TBARX, TBAX6, TBAX8],...|[key1, key2, key3,...|[[TBAX0, TBAXN, TBAXK],[TBAR7, TBAZJ, TBABI],...|
|          VRX8|{key39=[TBAX3, TBAX1, TBAX0], key46=[TBARQ, TBAX9, TBAX3],...|[key1, key2, key3,...|[[TBAX2, TBAXM, TBAXL],[TBAR8, TBAZK, TBABO],...|
|          CIZ2|{key39=[TBAX3, TBAXC, TBAX2], key46=[TBARA, TBAXQ, TBAX1],...|[key1, key2, key3,...|[[TBAX3, TBAXA, TBAXO],[TBAR9, TBAZL, TBABP],...|
|          BEO3|{key39=[TBAX9, TBAXQ, TBAX4], key46=[TBARP, TBAXV, TBAX2],...|[key1, key2, key3,...|[[TBAX1, TBAXS, TBAXI],[TBAR0, TBAZQ, TBABZ],...|
+--------------+-------------------------------------------------------------+---------------------+------------------------------------------------+

我尝试了以下方法,但出现以下错误:

getKeyValues = udf(lambda ar1, ar2: {ar1.get(x) for x in ar2})
df.withColumn("Values", getKeyValues(df["keyValue"], df["Keys"])).show()

AttributeError: 'unicode' object has no attribute 'get'

我在下面尝试了一个不同的 UDF,它也给出了相同的 unicode 错误:

def getKeyVals(ar1, ar2): 
    arr = []
    for x in ar2: 
        arr.append(ar1.get(x, None))
    return arr 
    
udf_split = udf(split, ArrayType(StringType()))

df.withColumn("test", udf_getKeyVals(df['keyValue'], df['Keys'])).show()

我也尝试了以下功能,但我遇到了类似的错误。 该函数来自以下MungingData webpage

def working_fun(mapping):
    def f(ar1):
        for x in ar1:
            return mapping.get(x)
    return F.udf(f)

df.withColumn("test", working_fun(df["keyValue"])(F.col('Keys'))).show()

如果有任何提示或建议,我们将不胜感激——谢谢!

更新后在下方包含架构和 Spark 版本详细信息

>>> df.select('reference_nbr', 'keyValue', 'Keys').schema.simpleString()
'struct<reference_nbr:string,keyValue:string,Keys:array<string>>'

>>> spark.version
u'2.3.2.3.1.5.6030-1'

【问题讨论】:

  • 请复制此命令df.select('reference_nbr', 'keyValue', 'Keys').schema.simpleString() 并将输出粘贴到您的问题中。另外,你的 spark 版本是什么?
  • 感谢@Kafels 的快速回复。我已经更新了原始帖子以包含架构和 Spark 版本信息。谢谢
  • 如果您不介意,您能否分享您的数据,并将此命令的输出改为df.select('reference_nbr', 'keyValue', 'Keys').limit(10).collect()
  • 嗨@Kafels 很抱歉回来晚了。不幸的是,“keyValue”和“Keys”列非常大,我认为我无法将它们发布到 Stack Overflow。如果我可以分享任何其他信息,请告诉我,再次感谢

标签: python dataframe apache-spark pyspark


【解决方案1】:

如果我们将keyvalue列从string类型转换为map类型,我们可以使用map_values函数来提取值:

我使用 UDF 将键值中的 = 替换为 : 以便我们可以使用 ast 模块将类型更改为映射。

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

spark = SparkSession.builder.master('local[*]').getOrCreate()


def get_dict(c):
    c = c.replace("=", ":")
    import ast
    dict_value = ast.literal_eval(c)
    return dict_value


get_dict_udf = udf(lambda c: get_dict(c), MapType(StringType(), ArrayType(IntegerType())))

# Sample dataframe
df = spark.createDataFrame(
    [('{"k3"= [6, 5, 4], "k1"= [4, 5, 1], "k8"= [8, 5, 6], "k5"= [7, 4, 3]}',
      ["k1", "k3", "k5", "k8"])]).toDF("keyvalue", "key")

df.withColumn("keyvalue", get_dict_udf("keyvalue")). \
    withColumn("values", sort_array(map_values("keyvalue"))).show(truncate=False)

+--------------------------------------------------------------------+----------------+--------------------------------------------+
|keyvalue                                                            |key             |values                                      |
+--------------------------------------------------------------------+----------------+--------------------------------------------+
|[k3 -> [6, 5, 4], k5 -> [7, 4, 3], k8 -> [8, 5, 6], k1 -> [4, 5, 1]]|[k1, k3, k5, k8]|[[4, 5, 1], [6, 5, 4], [7, 4, 3], [8, 5, 6]]|
+--------------------------------------------------------------------+----------------+--------------------------------------------+

【讨论】:

  • 感谢回复,谢谢!不幸的是,我收到以下错误:AttributeError: 'dict' object has no attribute 'replace'。我猜我可能不得不使用 regex_replace() 代替?
  • 好的,您可以使用,但在您的情况下,keyvalue 列是正确的字符串,因此 udf 中的变量 'c' 将是字符串类型,它应该接受 'c' 上的替换方法。
猜你喜欢
  • 2021-07-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-02-05
  • 1970-01-01
  • 1970-01-01
  • 2014-10-17
  • 1970-01-01
相关资源
最近更新 更多