核心是构建事件驱动流水线实现“切片—清洗—合并”三阶段解耦:分片需语义对齐(非字节均分),上传校验连续性,清洗由事件总线异步调度,合并通过原子写与redis状态同步保障一致性,并全程依托trace_id和上下文隔离实现错误恢复。

核心在于把“切片—清洗—合并”三阶段全部纳入事件驱动流水线,用异步非阻塞方式解耦各环节,避免线程/内存瓶颈,同时保证分片边界语义完整和结果可验证。
分片必须语义对齐,不能只按字节均分
直接用文件总长度除以线程数得到固定字节块,会截断 UTF-8 字符、切开 JSON 对象或 CSV 行,导致清洗逻辑出错。正确做法是:
- 每个分片起始位置由主线程预扫描确定:从 seek(offset) 开始,逐字节读直到遇到换行符(文本)、记录分隔符(如 \x00)、或 JSON object 边界(通过括号计数)
- 前端上传时携带 chunk_index 和 chunk_offset(字节级),服务端严格校验连续性,拒绝跳号、重复或越界分片
- 二进制格式(如 Parquet、ProtoBuf)可按固定块读,但需确保 chunk_offset + chunk_size ≤ file.length(),否则 read() 抛 EOFException
清洗任务交给异步事件总线调度
不在线程内同步执行清洗逻辑,而是将每个分片封装为事件,投递到事件总线,由独立消费者异步处理:
- 事件 payload 包含:file_id、chunk_index、offset、length、临时存储路径(如 S3 presigned URL 或本地 /tmp/chunk_123.bin)
- 消费者监听 topic(如 “file.chunk.ready”),拉取后加载数据 → 执行字段脱敏/编码转换/空值填充等清洗操作 → 输出清洗后 chunk 到指定位置
- 使用 Redis Streams 或 Kafka 作为事件总线,天然支持失败重试、消费位点回溯、多消费者负载均衡
合并阶段用原子写+状态同步保障一致性
所有清洗完成的分片需按序写入目标文件,且不能因并发写导致数据错位或残留:
- 合并前清空目标文件:new File(target).delete(); new File(target).createNewFile(); —— 避免 "rw" 模式下旧数据残留
- 每个分片写入使用 FileChannel.position(offset).write(buffer),而非复用 RandomAccessFile 实例;position() 调用是线程安全的,允许多个 channel 并发写不同 offset
- 引入轻量状态同步机制:用 Redis Hash 记录每个 chunk_index 的 status("cleaned" / "merged"),合并服务轮询确认全部分片就绪后再触发最终写入
上下文隔离与错误恢复要贯穿全程
单个大文件上传清洗过程可能持续数分钟,需防中断、防污染、防状态丢失:
- 每个上传请求绑定唯一 trace_id,并注入 Fiber Context(Swoole)或 StructuredTaskScope(Java 21),确保日志、指标、临时资源归属清晰
- 清洗失败时,事件总线自动重投(带退避),超过阈值则触发告警并标记该 chunk 为 "failed",后续合并跳过并记录缺失索引
- 系统重启后,通过扫描临时目录 + 查询 Redis 状态,自动续跑未完成的 chunk 清洗与合并任务











