如何基于SQL数据库关联查询结果生成多份JSON文件
实现方案与技巧
一、PySpark(Azure Databricks环境推荐)
1. 加载CSV数据
指定表头与数据类型加载三个表,避免自动推断引发的类型错误:
# 加载三个CSV表 pub_obj_df = spark.read.csv("/path/to/PublicationObject.csv", header=True, inferSchema=True) source_df = spark.read.csv("/path/to/Source.csv", header=True, inferSchema=True) source_details_df = spark.read.csv("/path/to/SourceObjectDetails.csv", header=True, inferSchema=True)
2. 执行多表关联
根据表间关联键(假设为PublicationObjectId、SourceId)完成关联,去除重复列:
# 关联Source与SourceObjectDetails source_with_details_df = source_df.join( source_details_df, source_df["SourceId"] == source_details_df["SourceId"], how="left" ).drop(source_details_df["SourceId"]) # 关联PublicationObject与上述结果 joined_df = pub_obj_df.join( source_with_details_df, pub_obj_df["PublicationObjectId"] == source_with_details_df["PublicationObjectId"], how="left" ).drop(source_with_details_df["PublicationObjectId"])
3. 构建嵌套JSON结构
通过分组聚合将关联数据打包为嵌套格式,匹配目标JSON结构:
from pyspark.sql import functions as F # 定义Source的嵌套结构 source_struct = F.struct( F.col("SourceId"), F.col("SourceName"), F.struct( F.col("DetailId"), F.col("DetailValue"), F.col("DetailType") ).alias("SourceObjectDetails") ) # 按PublicationObject分组,聚合关联的Source列表 nested_df = joined_df.groupBy( # 列出PublicationObject的所有字段,例如PublicationObjectId、Title、PublishDate等 "PublicationObjectId", "Title", "PublishDate" ).agg( F.collect_list(source_struct).alias("Sources") )
4. 输出单个JSON文件
针对500个PublicationObject的规模,采用逐行遍历写入的方式生成单PO单文件:
import json # 遍历每个PO数据,写入单独JSON文件 for row in nested_df.toLocalIterator(): po_id = row["PublicationObjectId"] # 注意:Azure Databricks中写入DBFS需加/dbfs/前缀 with open(f"/dbfs/path/to/output/{po_id}.json", "w") as f: json.dump(row.asDict(), f, indent=2)
二、Python Pandas方案(小数据量场景适用)
1. 加载并关联数据
import pandas as pd pub_obj_df = pd.read_csv("/path/to/PublicationObject.csv") source_df = pd.read_csv("/path/to/Source.csv") source_details_df = pd.read_csv("/path/to/SourceObjectDetails.csv") # 多表关联 source_with_details = pd.merge(source_df, source_details_df, on="SourceId", how="left") joined_df = pd.merge(pub_obj_df, source_with_details, on="PublicationObjectId", how="left")
2. 构建嵌套结构
# 分组生成嵌套数据 def build_nested_po(group): # 组装Sources列表 sources = [] for _, row in group.iterrows(): source_dict = { "SourceId": row["SourceId"], "SourceName": row["SourceName"], "SourceObjectDetails": { "DetailId": row["DetailId"], "DetailValue": row["DetailValue"], "DetailType": row["DetailType"] } } sources.append(source_dict) # 获取PO基础信息(取分组内第一行即可,同PO信息一致) po_base = group.iloc[0][["PublicationObjectId", "Title", "PublishDate"]].to_dict() po_base["Sources"] = sources return po_base # 分组处理所有PO nested_data = joined_df.groupby("PublicationObjectId").apply(build_nested_po).tolist()
3. 输出单个JSON文件
import json for po in nested_data: po_id = po["PublicationObjectId"] with open(f"/path/to/output/{po_id}.json", "w", encoding="utf-8") as f: json.dump(po, f, indent=2, ensure_ascii=False)
三、关键技巧
- PySpark性能优化:
- 提前定义Schema替代
inferSchema,减少类型推断开销,示例:from pyspark.sql.types import StructType, StructField, StringType, IntegerType pub_schema = StructType([ StructField("PublicationObjectId", IntegerType(), True), StructField("Title", StringType(), True), # 其他字段依次定义 ]) pub_obj_df = spark.read.csv("/path/to/file.csv", header=True, schema=pub_schema) - 使用
toLocalIterator()逐行处理数据,避免Driver内存溢出。
- 提前定义Schema替代
- JSON格式对齐:
- 用
alias()重命名字段,确保嵌套结构字段名与目标完全匹配;用fillna()处理空值,避免JSON中出现不必要的null。
- 用
- Azure Databricks路径注意:
- 本地文件操作访问DBFS需加
/dbfs/前缀,Spark API直接使用dbfs:/path/格式。
- 本地文件操作访问DBFS需加
- 验证环节:
- 先抽取10个以内的PO做测试,确认嵌套结构符合要求后再全量运行。
内容的提问来源于stack exchange,提问作者Saurabh Mehta
相关产品推荐
相关产品推荐

