analysis-task-2.md 4.5 KB

任务 2: 流式 ETL 监听器分析

StreamingEtlListener 核心逻辑分析

核心组件

  1. 数据清洗器 (dataCleaner)

    • 负责将原始数据转换为符合目标表结构的数据
    • 支持泛型设计,适应不同类型的源数据
  2. 数据加载器 (dataLoader)

    • 负责将清洗后的数据写入目标存储
    • 支持多种写入策略
  3. 脏数据处理器 (dirtyHandler)

    • 负责记录和存储清洗失败的数据
    • 确保错误数据不影响整体处理
  4. 计数器

    • totalRows:总行数计数器(线程安全)
    • sheetRows:当前工作表行数计数器(线程安全)
    • sheetFailRows:当前工作表失败行数计数器(线程安全)

核心方法分析

1. invoke(Map<Integer, String> rawData, AnalysisContext context)

执行流程:

  1. 初始化处理

    • 获取当前行索引 rowIndex
    • 增加总行数和当前工作表行数
    • 设置当前工作表信息 currentSheet
    • 跳过无工作表信息的情况
  2. 模板信息获取

    • 获取当前工作表的模板ID templateId
    • 从全局缓存获取模板对应的表配置信息 TableInfo
    • 获取清洗函数配置 funcRegx 和表头行号 lineNo
  3. 懒加载组件

    • 首次使用时创建数据加载器 dataLoader
    • 首次使用时创建数据清洗器 dataCleaner
  4. 数据分类处理

    • 表头行 (lineNo > rowIndex):收集表头字段到 headerList
    • 表头行 (lineNo == rowIndex):传递表头给清洗器进行字段映射准备
    • 数据行 (lineNo < rowIndex):执行清洗逻辑
  5. 数据清洗

    • 调用 dataCleaner.clean() 方法清洗数据
    • 处理清洗结果:
      • 数据有效:调用 dataLoader.write() 写入数据
      • 数据无效:记录错误信息,调用 dirtyHandler.handle() 处理脏数据
  6. 异常处理

    • 捕获处理过程中的异常
    • 记录错误日志,增加失败行数计数器

2. doAfterAllAnalysed(AnalysisContext context)

执行流程:

  1. 资源清理

    • 关闭数据加载器,刷新缓冲区
    • 刷新脏数据处理器,确保所有错误记录都被持久化
  2. 统计信息更新

    • 更新当前工作表的统计信息:
      • 数据行数 dataNum
      • 失败行数 failNum
      • 处理状态 FILE_CLEAN_COMPLETE
    • 将当前工作表添加到 extractSheetList
  3. 完成处理

    • 调用清洗器的 finish() 方法
  4. 状态重置

    • 重置所有状态,准备处理下一个工作表:
      • 清除 dataCleaner、dataLoader、currentSheet
      • 清空 headerList
      • 重置所有计数器
  5. 日志记录

    • 记录 ETL 完成信息,包括总处理行数

关键设计特点

  1. 流式处理:逐行处理数据,内存占用低,支持大数据量文件
  2. 懒加载:按需创建清洗器和加载器,节省资源
  3. 线程安全:使用 AtomicInteger 确保计数器的线程安全
  4. 错误处理:完善的错误捕获和处理机制,确保任务不会因单个数据行失败而中断
  5. 状态管理:清晰的状态重置机制,支持多工作表连续处理
  6. 扩展性:通过接口设计和策略模式,支持不同类型的数据处理

脏数据处理机制

  1. 错误捕获:捕获数据清洗过程中的异常
  2. 错误记录:创建 CleanErrorLog 对象,记录错误信息
  3. 错误处理:调用 dirtyHandler.handle() 方法处理脏数据
  4. 批量处理:在 doAfterAllAnalysed() 中调用 dirtyHandler.flush() 批量处理脏数据

数据流图

[原始数据行] → [invoke()方法]
               ↓
[设置当前工作表] → [获取模板信息]
               ↓
[懒加载组件] → [数据分类处理]
               ↓
[数据清洗] → [处理清洗结果]
               ↓
[写入数据/处理脏数据]

[所有数据处理完成] → [doAfterAllAnalysed()方法]
                   ↓
[资源清理] → [统计信息更新]
                   ↓
[完成处理] → [状态重置]

代码优化建议

  1. 错误处理增强:可以考虑添加更详细的错误分类和处理策略
  2. 性能优化:对于大量数据的情况,可以考虑批量处理机制
  3. 监控增强:添加处理速度、成功率等监控指标
  4. 配置灵活性:可以考虑将一些硬编码的配置参数化
  5. 异常恢复:添加异常恢复机制,确保部分失败不影响整体处理