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

PySpark中提取指定字段将JSON字符串转换为结构化DataFrame

PySpark处理Kafka Connect格式JSON数据

核心需求

从带有Kafka Connect schema结构的JSON中,提取payload下的c1、c2、create_ts、update_ts字段,并将毫秒级时间戳转换为YYYY-MM-dd HH:mm:ss格式的可读时间。

实现步骤

  1. 初始化SparkSession
    搭建基础运行环境:

    from pyspark.sql import SparkSession
    from pyspark.sql.functions import col, date_format, to_timestamp
    
    spark = SparkSession.builder \
        .appName("KafkaConnectJSONProcessing") \
        .getOrCreate()
    
  2. 读取原始JSON数据
    以读取本地JSON文件为例(若处理Kafka流数据,替换为readStream方式即可):

    # 读取JSON文件生成DataFrame
    df = spark.read.json("/path/to/your/kafka-connect-data.json")
    
  3. 提取payload层目标字段
    原始JSON分为schema和payload两层,直接展开payload提取所需列:

    payload_df = df.select(
        col("payload.c1"),
        col("payload.c2"),
        col("payload.create_ts"),
        col("payload.update_ts")
    )
    
  4. 转换毫秒时间戳为可读格式
    先将毫秒级时间戳转为Spark Timestamp类型,再格式化为指定字符串格式:

    processed_df = payload_df.withColumn(
        "create_ts",
        date_format(to_timestamp(col("create_ts") / 1000), "yyyy-MM-dd HH:mm:ss")
    ).withColumn(
        "update_ts",
        date_format(to_timestamp(col("update_ts") / 1000), "yyyy-MM-dd HH:mm:ss")
    )
    
  5. 查看最终结果
    输出结构化表格:

    processed_df.show(truncate=False)
    

    输出示例:

    +---+---+-------------------+-------------------+
    |c1 |c2 |create_ts          |update_ts          |
    +---+---+-------------------+-------------------+
    |67 |foo|2022-09-22 00:00:02|2022-09-22 00:00:02|
    +---+---+-------------------+-------------------+
    

补充说明

  • 若处理Kafka流数据,只需将read.json替换为readStream.format("kafka")的对应配置,转换逻辑完全通用。
  • c1字段为可选字段,若存在null值,Spark会自动保留,无需额外处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 17:01:26