使用AWS Glue处理含数组JSON:修正Schema识别与扁平化问题
解决AWS Glue爬虫无法正确推断动态字段JSON Schema的问题
我来帮你搞定这个问题——这种嵌套且字段不一致的JSON结构确实是Glue爬虫的常见痛点,尤其是当你想把data数组里的对象扁平化成单条记录的时候。下面是几个经过验证的解决方案,按可靠性和灵活性排序:
方案1:使用AWS Glue ETL脚本手动解析(最可靠)
因为data数组内的对象字段不固定,Glue爬虫自动推断Schema时会把所有文件里出现过的字段都列出来,没值的就填充null,这显然不是你想要的。手动写ETL脚本可以动态处理每个对象的字段,完美实现扁平化。
步骤分解:
- 在Glue Studio创建新的ETL作业,选择「Spark script editor」模式。
- 编写脚本读取S3中的原始JSON数据:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) # 替换成你的S3 JSON文件路径 df = spark.read.json("s3://your-bucket/path/to/json-files/")
- 展开
data数组,将每个数组元素转为单独的行:
from pyspark.sql.functions import explode # 展开data数组,每个元素作为单独的data_item exploded_df = df.select("device", "timestamp", explode("data").alias("data_item"))
- 动态提取
data_item中的所有字段,合并到主表:
from pyspark.sql.functions import col # 自动获取data_item里的所有字段(适配动态变化的字段) data_fields = exploded_df.select("data_item.*").columns # 将data_item的字段展开到主表,保留device和timestamp final_df = exploded_df.select("device", "timestamp", *[col(f"data_item.{field}").alias(field) for field in data_fields])
- 将处理后的数据写入S3(推荐用Parquet格式,优化Athena查询性能):
# 可以按device或timestamp分区,提升查询效率 final_df.write.partitionBy("device").parquet("s3://your-bucket/output-path/", mode="overwrite") # 可选:直接写入Glue数据目录,生成可查询的表 from awsglue.dynamicframe import DynamicFrame dynamic_df = DynamicFrame.fromDF(final_df, glueContext, "dynamic_df") glueContext.write_dynamic_frame.from_catalog( frame=dynamic_df, database="your-glue-database", table_name="your-target-table", transformation_ctx="write_to_catalog" ) job.commit()
运行作业后,Athena中查询目标表就能看到完全扁平化的单条记录,每个data对象的字段都会被保留,没有多余的null。
方案2:调整Glue爬虫+自定义分类器(适合字段变化有限的场景)
如果不想写ETL脚本,可以试试调整爬虫配置,配合精准的自定义分类器:
- 创建JSON分类器:指定JSON路径为
$(根节点),将data字段定义为array<map<string, string>>(用map代替固定struct,适配动态字段)。 - 配置爬虫时选择这个自定义分类器,关闭「Merge schemas」选项(默认开启会合并所有文件的字段,导致大量
null)。 - 爬虫完成后,在Athena中用
UNNEST展开数组,提取map中的键值对:
SELECT device, timestamp, -- 提取map中的字段,这里以第一个键值对为例,你可以根据需求扩展 map_keys(data_item)[0] AS field_name, map_values(data_item)[0] AS field_value FROM your-glue-table CROSS JOIN UNNEST(data) AS t(data_item)
这个方法的缺点是如果每个data对象有多个字段,需要手动处理每个键值对,灵活性不如ETL脚本。
方案3:预处理JSON文件(适合小数据集)
如果你的数据集不大,可以在上传到S3前预处理JSON,将data数组中的每个对象与device、timestamp合并为单独的JSON对象:
原始JSON示例:
{ "device": "device-001", "timestamp": "2024-05-20T14:30:00", "data": [ {"temperature": 26.5, "humidity": 58}, {"battery_level": 82} ] }
预处理后拆分两个独立JSON:
{"device": "device-001", "timestamp": "2024-05-20T14:30:00", "temperature": 26.5, "humidity": 58} {"device": "device-001", "timestamp": "2024-05-20T14:30:00", "battery_level": 82}
这样Glue爬虫就能直接推断出正确的Schema,Athena查询就是单条记录。可以用Python脚本批量处理:
import json import os input_dir = "/local/path/to/input/json" output_dir = "/local/path/to/output/json" os.makedirs(output_dir, exist_ok=True) for filename in os.listdir(input_dir): if filename.endswith(".json"): with open(os.path.join(input_dir, filename), "r") as f: raw_data = json.load(f) device = raw_data["device"] ts = raw_data["timestamp"] # 拆分data数组为单独对象 for idx, item in enumerate(raw_data["data"]): item["device"] = device item["timestamp"] = ts output_filename = f"{os.path.splitext(filename)[0]}_{idx}.json" with open(os.path.join(output_dir, output_filename), "w") as out_f: json.dump(item, out_f)
内容的提问来源于stack exchange,提问作者Maciej Malak
相关产品推荐
相关产品推荐

