【发布时间】: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