PySpark中提取指定字段将JSON字符串转换为结构化DataFrame
PySpark处理Kafka Connect格式JSON数据
核心需求
从带有Kafka Connect schema结构的JSON中,提取payload下的c1、c2、create_ts、update_ts字段,并将毫秒级时间戳转换为YYYY-MM-dd HH:mm:ss格式的可读时间。
实现步骤
初始化SparkSession
搭建基础运行环境:from pyspark.sql import SparkSession from pyspark.sql.functions import col, date_format, to_timestamp spark = SparkSession.builder \ .appName("KafkaConnectJSONProcessing") \ .getOrCreate()读取原始JSON数据
以读取本地JSON文件为例(若处理Kafka流数据,替换为readStream方式即可):# 读取JSON文件生成DataFrame df = spark.read.json("/path/to/your/kafka-connect-data.json")提取payload层目标字段
原始JSON分为schema和payload两层,直接展开payload提取所需列:payload_df = df.select( col("payload.c1"), col("payload.c2"), col("payload.create_ts"), col("payload.update_ts") )转换毫秒时间戳为可读格式
先将毫秒级时间戳转为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") )查看最终结果
输出结构化表格: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ك
相关产品推荐
相关产品推荐

