如何在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
相关产品推荐
相关产品推荐

