在Azure Databricks笔记本中将PySpark Row转为JSON并保留字段名
问题描述
在Azure Databricks笔记本中有如下PySpark Row类型数据:
indv_msg = [Row(cbm_json_output=Row(country_code='USA', date='06-10-2023', date_epoch='1696550400', id='USA-001535-1696550400', interfaceVersion='1.0.0', opmode_car_door=Row(health_category='GREEN', msg_id='1', num_yellow_preds_in_last_14_days=0, reason=None, reasonDetail=None), opmode_landing_door=Row(health_category='GREEN', msg_id='1', reason=None, reasonDetail=None), sensor=Row(component_type=None, health_category=None, landing_priority=None, msg_id='1', num_yellow_preds_in_last_14_days=None, reason=None, reasonDetail=None), unit_id='001535'))]
使用json.dumps(indv_msg, indent=2)转换为JSON字符串时,country_code、date等字段名被省略,仅输出对应值,转换结果仅保留数值列表。需要得到保留键值对格式的JSON字符串,比如包含"country_code":"USA"这类键值映射的结构。
解决方案
PySpark的Row对象本质是带字段名的元组,直接用json.dumps序列化时会默认按元组处理,只输出值而丢失键名。要保留键值对,需先把嵌套的Row结构转换为Python字典,再进行JSON序列化。
方法1:递归转换Row为字典
编写递归函数,将所有嵌套的Row对象转为字典:
import json from pyspark.sql import Row def row_to_dict(row): if isinstance(row, Row): return {k: row_to_dict(v) for k, v in row.asDict().items()} elif isinstance(row, list): return [row_to_dict(item) for item in row] else: return row # 转换并序列化 dict_msg = row_to_dict(indv_msg) json_str = json.dumps(dict_msg, indent=2) print(json_str)
方法2:Spark内置函数批量处理(适用于DataFrame场景)
如果是处理DataFrame中的批量数据,可直接使用Spark的to_json函数,自动保留键值对结构:
from pyspark.sql.functions import to_json, struct # 假设df包含cbm_json_output列 df = df.withColumn("json_output", to_json(struct("cbm_json_output.*")))
以上两种方法都能生成包含完整键值对的JSON字符串,满足需求。
内容的提问来源于stack exchange,提问作者Saswat Ray
相关产品推荐
相关产品推荐

