1. 为什么需要流式处理大数据?

想象一下你要从水库里取水浇灌农田。如果一次性把水库抽干,不仅需要巨大的储水容器,还可能压垮运输管道。更合理的做法是用水管持续引流,按需取用——这就是流式处理的核心逻辑。

在数据处理领域,传统的一次性加载方式(如List<User> users = userMapper.selectAll())在处理百万级数据时,相当于把整个水库的水瞬间灌进内存。我曾在生产环境遇到过因此导致的内存溢出(OOM),JVM堆内存直接飙到8GB上限,服务崩溃的惨痛教训。

流式处理的优势在于:

  • 内存可控:每次只处理单条或小批量数据
  • 实时响应:数据边读取边处理,无需等待全量加载
  • 资源节约:减少数据库连接占用时间

2. MyBatis流式处理双剑客:ResultHandler与游标

2.1 ResultHandler的工作原理

ResultHandler是MyBatis提供的回调接口,其核心方法是:

void handleResult(ResultContext<?> context);

当数据库返回结果集时,MyBatis会逐行调用这个方法。我常用它来做这些事:

  • 实时写入文件(如CSV导出)
  • 分批插入到其他数据库
  • 内存计算(如实时统计)

这里有个实战技巧:通过ResultContext对象可以获取当前处理状态:

// 获取当前行数据
User user = (User)context.getResultObject(); 
// 判断是否停止处理
if(shouldStop()) context.stop();

2.2 数据库游标的关键配置

单纯使用ResultHandler还不够,必须配合正确的游标类型。重点配置这两个参数:

@Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 1000)
void streamData(ResultHandler<User> handler);
  • FORWARD_ONLY:声明只进游标,禁止随机访问(节省内存)
  • fetchSize:控制每次从数据库拉取的行数(实测1000-5000是较优区间)

踩坑提醒:MySQL驱动默认会缓存所有结果集到内存,必须同时在JDBC连接字符串添加useCursorFetch=true参数才能真正启用流式:

jdbc:mysql://host/db?useCursorFetch=true

3. 构建完整流式处理管道

3.1 基础实现方案

完整代码示例:

// 1. 自定义处理器
public class UserExportHandler implements ResultHandler<User> {
    private final CSVWriter writer;
    
    @Override
    public void handleResult(ResultContext<? extends User> ctx) {
        User user = ctx.getResultObject();
        writer.writeRow(userToArray(user));
        if(ctx.getResultCount() % 1000 == 0) {
            writer.flush(); // 每1000行刷盘
        }
    }
}

// 2. Mapper接口配置
@Options(resultSetType = FORWARD_ONLY, fetchSize = 1000)
void streamUsers(@Param("minId") long minId, ResultHandler<User> handler);

// 3. 调用示例
try(CSVWriter writer = new CSVWriter(new FileWriter("output.csv"))) {
    UserExportHandler handler = new UserExportHandler(writer);
    userMapper.streamUsers(1000L, handler);
}

3.2 高并发优化方案

当需要并行处理时,推荐采用生产者-消费者模式:

// 创建有界队列防止内存膨胀
BlockingQueue<User> queue = new ArrayBlockingQueue<>(1000); 

// ResultHandler作为生产者
public class QueueProducer implements ResultHandler<User> {
    @Override
    public void handleResult(ResultContext<? extends User> ctx) {
        queue.put(ctx.getResultObject()); // 阻塞式插入
    }
}

// 启动消费者线程池
ExecutorService pool = Executors.newFixedThreadPool(4);
for(int i=0; i<4; i++) {
    pool.submit(() -> {
        while(true) {
            User user = queue.take();
            processUser(user);
        }
    });
}

性能数据:在16核服务器上测试处理1000万条数据:

  • 单线程耗时:182秒
  • 4消费者线程:63秒
  • 内存占用始终稳定在500MB以内

4. 生产环境实战经验

4.1 事务管理要点

流式处理中常见的事务误区:

// 错误示例:整个流式过程在事务中
@Transactional 
public void exportData() {
    mapper.streamData(handler); // 可能导致连接长时间占用
}

正确做法应该是:

  1. 分批次开启事务(每1000条提交一次)
  2. 使用TransactionTemplate编程式事务
  3. 对于只读操作,设置@Transactional(readOnly=true)

4.2 资源释放技巧

必须确保流关闭后释放资源。推荐try-with-resources模式:

try(SqlSession session = sqlSessionFactory.openSession()) {
    UserMapper mapper = session.getMapper(UserMapper.class);
    mapper.streamData(new MyHandler());
} // 自动关闭session和结果集

4.3 异常处理方案

流式处理中异常需要特殊处理:

public class SafeResultHandler implements ResultHandler<User> {
    @Override
    public void handleResult(ResultContext<? extends User> ctx) {
        try {
            businessLogic(ctx.getResultObject());
        } catch (Exception e) {
            ctx.stop(); // 终止流
            throw new RuntimeException("处理失败", e);
        }
    }
}

5. 性能对比与调优

5.1 不同方案内存占用对比

处理方式100万条数据内存峰值耗时
传统List查询1.2GB28s
基础流式处理50MB32s
并行流式处理80MB11s

5.2 关键参数调优

  1. fetchSize

    • MySQL建议1000-5000
    • Oracle建议100-200
    • PostgreSQL建议500-2000
  2. JVM参数

    -XX:+UseG1GC -Xmx2g -Xms2g 
    

    使用G1垃圾回收器能更好处理大量短期对象

  3. 批处理大小: 写入数据库时,每500-1000条提交一次事务效率最高

6. 典型应用场景

6.1 数据迁移工具

最近用这套方案实现了跨库数据迁移:

public class DataMigrator implements ResultHandler<SourceData> {
    private final TargetMapper targetMapper;
    private int batchCount = 0;
    
    @Override
    public void handleResult(ResultContext<? extends SourceData> ctx) {
        targetMapper.insert(convert(ctx.getResultObject()));
        if(++batchCount % 1000 == 0) {
            sqlSession.commit(); // 分批提交
        }
    }
}

6.2 实时报表生成

金融行业日终报表处理示例:

public class ReportGenerator {
    public void generateDailyReport(LocalDate date) {
        ReportHandler handler = new ReportHandler();
        // 流式处理交易数据
        transactionMapper.streamByDate(date, handler);
        // 实时生成统计结果
        handler.saveReport();
    }
}

6.3 大数据预处理

配合Spark等大数据框架时,可以这样对接:

JavaRDD<User> rdd = sparkContext.emptyRDD();
ResultHandler handler = context -> {
    rdd.union(sparkContext.parallelize(
        Arrays.asList(context.getResultObject())));
};
userMapper.streamAll(handler);

更多推荐