根据您的示例,您可以使用 SparkSQL 函数 str_to_map 将字符串转换为 Map,然后从所需的映射键中选择值(以下代码假设 StringType 列名称为 value):
from pyspark.sql import functions as F
keys = ['Type', 'Model', 'ID', 'conn seq']
df.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \
.select("*", *[ F.col('m')[k].alias(k) for k in keys ]) \
.show()
+--------------------+--------------------+---------+-----+---+--------+
| value| m| Type|Model| ID|conn seq|
+--------------------+--------------------+---------+-----+---+--------+
|Type=<Series VR> ...|[Type -> Series V...|Series VR| 1Ac4| 34| 2|
|Type=<SeriesX> Mo...|[Type -> SeriesX,...| SeriesX| 12Q3|231| 3423123|
+--------------------+--------------------+---------+-----+---+--------+
注意:这里我们使用正则表达式模式> * 来分割对,模式=< 来分割键/值。检查this link 如果映射的keys 是动态的并且无法预定义,请确保过滤掉EMPTY 键。
编辑:基于 cmets,对映射键进行不区分大小写的搜索。对于 Spark 2.3,我们可以在使用 str_to_map 函数之前使用pandas_udf 预处理value 列:
-
为匹配的键设置正则表达式模式(在捕获组 1 中)。这里我们使用(?i)来设置不区分大小写的匹配,并添加两个锚点\b和(?==),这样匹配的子串必须在左边有一个单词边界,后面跟着一个=标记右边。
ptn = "(?i)\\b({})(?==)".format('|'.join(keys))
print(ptn)
#(?i)\b(Type|Model|ID|conn seq)(?==)
-
设置 pandas_udf 以便我们可以使用 Series.str.replace() 并设置一个回调(小写 $1)作为替换:
lower_keys = F.pandas_udf(lambda s: s.str.replace(ptn, lambda m: m.group(1).lower()), "string")
-
将所有匹配的键转换为小写:
df1 = df.withColumn('value', lower_keys('value'))
+-------------------------------------------------------+
|value |
+-------------------------------------------------------+
|type=<Series VR> model=<1Ac4> id=<34> conn seq=<2> |
|type=<SeriesX> model=<12Q3> id=<231> conn seq=<3423123>|
+-------------------------------------------------------+
-
使用str_to_map创建map,然后以k.lower()为key,查找对应的值。
df1.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \
.select("*", *[ F.col('m')[k.lower()].alias(k) for k in keys ]) \
.show()
注意:如果您以后可以使用Spark 3.0+,请跳过上述步骤并改用transform_keys函数:
df.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \
.withColumn("m", F.expr("transform_keys(m, (k,v) -> lower(k))")) \
.select("*", *[ F.col('m')[k.lower()].alias(k) for k in keys ]) \
.show()
对于 Spark 2.4+,将 transform_keys(...) 替换为以下内容:
map_from_entries(transform(map_keys(m), k -> (lower(k), m[k])))