You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中从字符串类型DataFrame列提取键值对

解决PySpark提取非标准键值对字符串字段的问题

你的address_info列不是标准JSON格式(用=替代了:,且无引号),所以直接用from_json会失败。这里提供两种可行方案:

方案一:转换为标准JSON后解析

先将非标准字符串转为标准JSON格式,再用from_json解析提取字段:

步骤与代码

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType
from pyspark.sql.functions import regexp_replace, from_json, col

# 初始化SparkSession
spark = SparkSession.builder.appName("ExtractAddressData").getOrCreate()

# 测试数据
data = [("XXX", "{original_email=test@sparp.org, fqdn=sparp.org, domain=sharp, host=sharp.org, subdomain=, alias=, addr=test, tld=org, normalized_email=test@sharp.org}")]
df = spark.createDataFrame(data, ["name", "address_info"])

# 定义需要提取字段的Schema
address_schema = StructType([
    StructField("original_email", StringType(), nullable=True),
    StructField("fqdn", StringType(), nullable=True),
    StructField("domain", StringType(), nullable=True)
])

# 将非标准字符串转换为标准JSON
df_with_json = df.withColumn(
    "standard_json",
    regexp_replace(
        regexp_replace(
            regexp_replace(col("address_info"), "^\\{|\\}$", ""),  # 移除首尾大括号
            "([^,=]+)=([^,]+)", "\"$1\":\"$2\""),  # 替换key=value为"key":"value"
        "=,", "=\"\","  # 处理空值(如subdomain=, 转为"subdomain":"",)
    )
).withColumn("standard_json", regexp_replace(col("standard_json"), "=$", "\""))  # 处理末尾空值

# 解析JSON并提取目标字段
result_df = df_with_json.withColumn(
    "address_struct", from_json(col("standard_json"), address_schema)
).select(
    "name",
    col("address_struct.original_email").alias("original_email"),
    col("address_struct.fqdn").alias("fqdn"),
    col("address_struct.domain").alias("domain")
)

# 查看结果
result_df.show()

方案二:直接正则提取目标字段

如果只需要提取特定几个字段,用regexp_extract直接匹配键值对,无需转换整个JSON,效率更高:

代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import regexp_extract, col

spark = SparkSession.builder.appName("ExtractAddressData").getOrCreate()

# 测试数据
data = [("XXX", "{original_email=test@sparp.org, fqdn=sparp.org, domain=sharp, host=sharp.org, subdomain=, alias=, addr=test, tld=org, normalized_email=test@sharp.org}")]
df = spark.createDataFrame(data, ["name", "address_info"])

# 直接提取指定字段
result_df = df.select(
    col("name"),
    regexp_extract(col("address_info"), r"original_email=([^,]+)", 1).alias("original_email"),
    regexp_extract(col("address_info"), r"fqdn=([^,]+)", 1).alias("fqdn"),
    regexp_extract(col("address_info"), r"domain=([^,]+)", 1).alias("domain")
)

result_df.show()

正则说明

r"key=([^,]+)" 匹配key=后直到下一个逗号的所有内容,捕获组1即为对应字段的值,空值会自动提取为空字符串。

内容的提问来源于stack exchange,提问作者Fkp

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 09:00:20