You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将含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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 01:24:28