如何将含Job/Stage的Spark对账报告文本转为层级JSON文件
解决方案
1. 明确目标JSON结构
首先要定义好最终的层级结构,确保所有Job归到Jobs数组,每个Job下的Stage归到自身的Stages数组,环境信息单独放在SparkEnv节点:
{ "SparkEnv": { "spark.version": "3.3.0", "spark.executor.cores": "4" // 其他环境属性 }, "Jobs": [ { "JobID": "1", "JobName": "UserDataImport", "Duration": "5min", "Stages": [ { "StageID": "1.1", "TaskCount": "100", "SuccessRate": "100%" }, { "StageID": "1.2", "TaskCount": "50", "SuccessRate": "98%" } ] } // 其他Job对象 ] }
2. 核心解析逻辑
通过逐行扫描+状态标记的方式,区分当前处理的是环境信息、Job属性还是Stage属性:
- 初始化结果字典,包含
SparkEnv(空字典)和Jobs(空数组)。 - 用变量标记当前所处的层级:
current_job(指向当前正在处理的Job对象)、current_stage(指向当前正在处理的Stage对象)、in_env_section(是否在环境信息区块)。 - 逐行处理文本:
- 遇到环境信息的起始标记,切换到环境层级。
- 遇到新Job的标记,创建新的Job对象并加入
Jobs数组,设置current_job为这个对象,同时初始化Stages数组。 - 遇到新Stage的标记,创建新的Stage对象并加入当前Job的
Stages数组,设置current_stage为这个对象。 - 遇到键值对行,根据当前层级将键值对存入对应节点。
3. 代码实现(Python)
假设你的对账报告文本格式是类似分段的键值对(如示例文本),以下是可直接复用的解析代码:
import re import json def parse_spark_reconciliation_report(file_path): # 初始化目标结构 result = {"SparkEnv": {}, "Jobs": []} current_job = None current_stage = None in_env_section = False with open(file_path, 'r', encoding='utf-8') as f: for line in f: line = line.strip() # 跳过空行 if not line: continue # 切换到环境信息区块 if line.startswith("Spark Environment Information:"): in_env_section = True continue # 匹配新Job(格式示例:Job ID: 1) job_match = re.match(r"Job ID: (\S+)", line) if job_match: in_env_section = False # 创建新Job对象 current_job = { "JobID": job_match.group(1), "Stages": [] } result["Jobs"].append(current_job) current_stage = None continue # 匹配新Stage(格式示例:Stage ID: 1.1) stage_match = re.match(r"Stage ID: (\S+)", line) if stage_match and current_job: current_stage = { "StageID": stage_match.group(1) } current_job["Stages"].append(current_stage) continue # 处理键值对(格式示例:spark.version: 3.3.0) if ":" in line: key, value = line.split(":", 1) key = key.strip() value = value.strip() if in_env_section: result["SparkEnv"][key] = value elif current_stage: current_stage[key] = value elif current_job: current_job[key] = value return result # 解析并生成JSON文件 parsed_data = parse_spark_reconciliation_report("spark_report.txt") with open("spark_report_structured.json", 'w', encoding='utf-8') as f: json.dump(parsed_data, f, indent=2, ensure_ascii=False)
4. 适配不同报告格式的调整要点
- 如果报告用分隔线(如
=== Job 1 ===)区分Job,修改正则表达式为re.match(r"=== Job (\S+) ===", line)。 - 如果属性值存在换行(比如多行描述),需要添加逻辑判断下一行是否属于当前字段的值,暂存到临时变量中。
- 如果存在重复键,可根据需求选择覆盖或转为数组存储。
内容的提问来源于stack exchange,提问作者anonymous
相关产品推荐
相关产品推荐

