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

求基于条件在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 13:45:12