咨询基于嵌套JSON构建Athena表或仪表盘的替代方案
方案1:Athena嵌套类型+UNNEST直接查询(最快上手)
你的JSON里Job和Stage都是动态命名的键(Job_ID 0、Job_ID 1...),这种结构用Glue爬虫容易报错,但Athena支持直接定义嵌套的MAP类型适配,不用提前扁平化。
第一步:创建嵌套结构的Athena表
先写DDL定义表结构,把动态的Job和Stage字段定义为MAP类型,每个MAP的值是对应的结构体:
CREATE EXTERNAL TABLE IF NOT EXISTS emr_spark_reports ( report_for string, execution_of string, spark_version string, java_version string, start_time string, jobs map<string, struct< submission_time: string, run_time: string, result: string, number_of_stages: string, stages: map<string, struct< number_of_tasks: string, number_of_executed_tasks: string, completion_time: string, stage_executed: string >> >> ) ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe' LOCATION 's3://你的存储桶路径/存放JSON的文件夹/' TBLPROPERTIES ( 'serialization.format' = '1' );
注意:如果你的JSON是单个大文件(不是每行一个JSON对象),需要在TBLPROPERTIES里加上
'has_encrypted_data'='false',确保SerDe能读取整个文件的JSON结构。
第二步:用UNNEST展开嵌套数据
建表后,直接用UNNEST把Job和Stage拆成扁平化的行,方便分析:
SELECT report_for, execution_of, spark_version, java_version, start_time, job_key AS job_id, job_val.submission_time, job_val.run_time, job_val.result, job_val.number_of_stages, stage_key AS stage_id, stage_val.number_of_tasks, stage_val.number_of_executed_tasks, stage_val.completion_time, stage_val.stage_executed FROM emr_spark_reports CROSS JOIN UNNEST(jobs) AS t(job_key, job_val) CROSS JOIN UNNEST(job_val.stages) AS s(stage_key, stage_val);
可以把这个查询保存为视图,后续直接用视图做分析或对接仪表盘。
方案2:用Athena CTAS生成扁平化表(适合长期使用)
如果需要频繁查询,直接用CTAS把上面的查询结果存成Parquet表,性能更好:
CREATE TABLE emr_spark_flattened WITH ( format = 'PARQUET', external_location = 's3://你的存储桶路径/扁平化表的存储路径/', partitioned_by = ARRAY['execution_of'] -- 按应用ID分区,优化查询速度 ) AS SELECT report_for, execution_of, spark_version, java_version, start_time, regexp_extract(job_key, 'Job_ID (\d+)', 1) AS job_id, -- 提取纯数字Job ID job_val.submission_time, regexp_extract(job_val.run_time, '(\d+)', 1)::int AS run_time_ms, -- 转换为数值型 job_val.result, job_val.number_of_stages::int AS number_of_stages, regexp_extract(stage_key, 'Stage_ID (\d+)', 1) AS stage_id, -- 提取纯数字Stage ID stage_val.number_of_tasks::int AS number_of_tasks, stage_val.number_of_executed_tasks::int AS number_of_executed_tasks, regexp_extract(stage_val.completion_time, '(\d+)', 1)::int AS completion_time_ms, stage_val.stage_executed FROM emr_spark_reports CROSS JOIN UNNEST(jobs) AS t(job_key, job_val) CROSS JOIN UNNEST(job_val.stages) AS s(stage_key, stage_val);
生成的Parquet表可以直接用于所有后续查询,性能比直接查JSON提升明显。
方案3:直接对接QuickSight做仪表盘
不管用上面的视图还是CTAS表,都可以直接把Athena作为数据源连接到QuickSight:
- 在QuickSight里添加Athena数据源,选择对应的视图或表
- 拖拽字段创建可视化:比如按Job ID统计运行时间、按Stage查看任务执行率、按Result统计作业成功率等
- 可以在QuickSight里把时间字段(
start_time、submission_time)转换成日期类型,做时间趋势分析
关键注意点
- 确保Athena有读取S3源文件、写入查询结果/CTAS表的权限
- 如果需要对运行时间等数值做计算,用
regexp_extract提取数字后转换为数值类型,方便后续统计 - 动态键的纯数字ID可以用正则提取,让结果更整洁
内容的提问来源于stack exchange,提问作者anonymous
相关产品推荐
相关产品推荐

