求基于条件在PySpark中拆分Status_MP_c列生成sent/open/click列的逻辑
PySpark 按条件拆分状态列生成新列
核心思路
从Status_MP_c列中提取sent、open、click三个状态值,当对应状态为true时,新列保留该布尔值(或转换为1),否则设为false(或0)。以下是针对不同数据格式的实现方案:
方案1:Status_MP_c为结构体类型(StructType)
如果源数据中Status_MP_c本身已经是包含sent、open、click字段的结构体,直接提取字段即可:
from pyspark.sql import SparkSession # 初始化SparkSession(若已存在可跳过) spark = SparkSession.builder.appName("SplitStatus").getOrCreate() # 读取源文件(示例为CSV,根据实际格式调整) df = spark.read.csv("your_file_path.csv", header=True, inferSchema=True) # 拆分状态列生成新列 df_result = df.withColumn("sent", df["Status_MP_c.sent"]) \ .withColumn("open", df["Status_MP_c.open"]) \ .withColumn("click", df["Status_MP_c.click"]) # 查看结果 df_result.show(truncate=False)
方案2:Status_MP_c为JSON字符串格式
若Status_MP_c是JSON格式的字符串(如{"sent":true,"open":false,"click":true}),先解析为结构体再提取字段:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json from pyspark.sql.types import StructType, StructField, BooleanType spark = SparkSession.builder.appName("ParseStatusJSON").getOrCreate() df = spark.read.csv("your_file_path.csv", header=True, inferSchema=True) # 定义JSON对应的Schema status_schema = StructType([ StructField("sent", BooleanType(), nullable=True), StructField("open", BooleanType(), nullable=True), StructField("click", BooleanType(), nullable=True) ]) # 解析JSON并生成新列 df_result = df.withColumn("status_struct", from_json(col("Status_MP_c"), status_schema)) \ .withColumn("sent", col("status_struct.sent")) \ .withColumn("open", col("status_struct.open")) \ .withColumn("click", col("status_struct.click")) \ .drop("status_struct") # 可选:删除中间临时列 df_result.show(truncate=False)
方案3:Status_MP_c为自定义分隔格式
如果Status_MP_c是类似sent=true;open=false;click=true的分隔字符串,用正则提取状态值:
from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract from pyspark.sql.types import BooleanType spark = SparkSession.builder.appName("ExtractStatus").getOrCreate() df = spark.read.csv("your_file_path.csv", header=True, inferSchema=True) # 用正则提取每个状态的布尔值 df_result = df.withColumn("sent", regexp_extract(col("Status_MP_c"), r"sent:(true|false)", 1).cast(BooleanType())) \ .withColumn("open", regexp_extract(col("Status_MP_c"), r"open:(true|false)", 1).cast(BooleanType())) \ .withColumn("click", regexp_extract(col("Status_MP_c"), r"click:(true|false)", 1).cast(BooleanType())) df_result.show(truncate=False)
可选:将布尔值转换为0/1
如果需要把true/false转为1/0用于统计,可使用when函数:
from pyspark.sql.functions import when df_result = df_result.withColumn("sent", when(col("sent") == True, 1).otherwise(0)) \ .withColumn("open", when(col("open") == True, 1).otherwise(0)) \ .withColumn("click", when(col("click") == True, 1).otherwise(0))
内容的提问来源于stack exchange,提问作者Athiya
相关产品推荐
相关产品推荐

