# 任务 2: 流式 ETL 监听器分析 ## StreamingEtlListener 核心逻辑分析 ### 核心组件 1. **数据清洗器** (`dataCleaner`) - 负责将原始数据转换为符合目标表结构的数据 - 支持泛型设计,适应不同类型的源数据 2. **数据加载器** (`dataLoader`) - 负责将清洗后的数据写入目标存储 - 支持多种写入策略 3. **脏数据处理器** (`dirtyHandler`) - 负责记录和存储清洗失败的数据 - 确保错误数据不影响整体处理 4. **计数器** - `totalRows`:总行数计数器(线程安全) - `sheetRows`:当前工作表行数计数器(线程安全) - `sheetFailRows`:当前工作表失败行数计数器(线程安全) ### 核心方法分析 #### 1. `invoke(Map 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. **异常恢复**:添加异常恢复机制,确保部分失败不影响整体处理