PySpark中拆分pur_details字符串列提取指定字段(适配百万级数据)
高效提取PySpark字符串列中的指定字段(百万级数据适配)
原始数据表
+------+---------------------------------+---------------+-------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |cus_id|cus_nm |pur_region |purchase_dt |pur_details | +------+---------------------------------+---------------+-------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |0121 |Johnny |USA |2023-01-12 |[{product_id=XA8096521JKAZ42F123, product_name=luxury_watch_collection_rolex_GZ, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0137 |Kevin J Brown |USA |2022-05-31 |[{product_id=XA14567JKR700135126, product_name=luxury_watch_collection_rolex_LA, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0168 |Patrikson |UK |2022-11-08 |[{product_id=XAHJYZK906423623571, product_name=luxury_watch_collection_gucci_09, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0365 |Ryan Ray |USA |2021-10-12 |[{product_id=XAOPLKR7520HJV00109, product_name=luxury_watch_collection_vancleef, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |2600 |Jay |AUS |2022-11-11 |[{product_id=XA096534987GGHJLRAC, product_name=sports_eyewear, description=athlete sports sun glasses, check=sale_item, sale_price_gap=BOGO 20% off, sale_vendor=mrporter.com, action=report}] | +------+---------------------------------+---------------+-------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
表结构Schema
root |-- cus_id: string (nullable = true) |-- cus_nm: string (nullable = true) |-- pur_region: string (nullable = true) |-- purchase_dt: string (nullable = true) |-- pur_details: string (nullable = true)
需求说明
从pur_details字符串列中提取check和sale_price_gap作为独立列,若目标字段不存在,对应列值设为null;方案需适配百万级数据量,保证执行效率。
预期输出
+------+---------------------------------+---------------+-------------+----------+---------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |cus_id|cus_nm |pur_region |purchase_dt |check |sale_price_gap |pur_details | +------+---------------------------------+---------------+-------------+----------+---------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |0121 |Johnny |USA |2023-01-12 |sale_item |upto 30% on_sale |[{product_id=XA8096521JKAZ42F123, product_name=luxury_watch_collection_rolex_GZ, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0137 |Kevin J Brown |USA |2022-05-31 |sale_item |upto 30% on_sale |[{product_id=XA14567JKR700135126, product_name=luxury_watch_collection_rolex_LA, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0168 |Patrikson |UK |2022-11-08 |sale_item |upto 30% on_sale |[{product_id=XAHJYZK906423623571, product_name=luxury_watch_collection_gucci_09, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |0365 |Ryan Ray |USA |2021-10-12 |sale_item |upto 30% on_sale |[{product_id=XAOPLKR7520HJV00109, product_name=luxury_watch_collection_vancleef, description=mens watch round dail on sale, check=sale_item, tag=watch, sale_price_gap=upto 30% on_sale, sale_vendor=mrporter.com, action=entry}] | |2600 |Jay |AUS |2022-11-11 |sale_item |BOGO 20% off |[{product_id=XA096534987GGHJLRAC, product_name=sports_eyewear, description=athlete sports sun glasses, check=sale_item, sale_price_gap=BOGO 20% off, sale_vendor=mrporter.com, action=report}] | +------+---------------------------------+---------------+-------------+----------+---------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
高效实现方案
核心思路
优先使用PySpark内置函数而非自定义UDF,避免Python-JVM序列化开销,适配大数据量处理:
- 将非标准的
key=value格式字符串转换为标准JSON格式; - 用
from_json解析JSON字符串为Struct类型; - 直接提取目标字段,缺失时自动返回
null。
代码实现
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("ExtractPurDetails").getOrCreate() # 加载原始数据(替换为你的实际数据加载逻辑) # df = spark.read.table("your_table_name") # 定义解析pur_details的Schema pur_detail_schema = StructType([ StructField("check", StringType(), nullable=True), StructField("sale_price_gap", StringType(), nullable=True), ]) # 转换非标准格式为标准JSON df_transformed = df.withColumn( "pur_details_json", regexp_replace( regexp_replace( regexp_replace(col("pur_details"), r"^\[{|}\]$", ""), # 移除首尾的[{和}] r"(\w+)=([^,]+)", r'"\1":"\2"' # 将key=value转为"key":"value"格式 ), r",\s*", r"," # 清理逗号后的空格,保证JSON格式严谨 ) ) # 解析JSON并提取目标字段 result_df = df_transformed.withColumn( "pur_details_struct", from_json(col("pur_details_json"), pur_detail_schema) ).select( "cus_id", "cus_nm", "pur_region", "purchase_dt", col("pur_details_struct.check").alias("check"), col("pur_details_struct.sale_price_gap").alias("sale_price_gap"), "pur_details" ).drop("pur_details_json", "pur_details_struct") # 查看结果 result_df.show(truncate=False)
效率说明
- 所有操作均为Spark内置的矢量化操作,避免了UDF的逐行处理开销,适合百万级及以上数据量;
- 正则替换仅处理必要的格式转换,逻辑简洁高效;
from_json由Spark引擎优化执行,性能远优于自定义解析逻辑。
内容的提问来源于stack exchange,提问作者1ksj8jdnu36flksf
相关产品推荐
相关产品推荐

