【发布时间】:2017-07-18 07:32:58
【问题描述】:
我有一个包含 3+ 百万条记录的表。我需要从数据库中读取所有这些记录,将它们发送到 kafka 队列以供其他系统处理。然后从输出 kafka 队列中读取结果并写回数据库。
我需要在理智的部分进行读写,否则我会立即收到 OOM 异常。
用mybatis实现批量读写操作有哪些可能的技术方案?
非常感谢整洁的工作示例。
【问题讨论】:
标签: java mybatis spring-mybatis
我有一个包含 3+ 百万条记录的表。我需要从数据库中读取所有这些记录,将它们发送到 kafka 队列以供其他系统处理。然后从输出 kafka 队列中读取结果并写回数据库。
我需要在理智的部分进行读写,否则我会立即收到 OOM 异常。
用mybatis实现批量读写操作有哪些可能的技术方案?
非常感谢整洁的工作示例。
【问题讨论】:
标签: java mybatis spring-mybatis
我会写伪代码,因为我对 Kafka 不太了解。
首先在读取时,Mybatis 默认行为是以 List 的形式返回结果,但你不想将 300 万个对象加载到内存中。这必须通过使用org.apache.ibatis.session.ResultHandler<T>的自定义实现来覆盖
public void handleResult(final ResultContext<YourType> context) {
addToKafkaQueue(context.getResultObject());
}
如果Mybatis全局设置中没有定义值,也要设置语句的fetchSize(当使用基于注解的映射器时:@Option(fetchSize=500))。如果未设置,此选项默认依赖于驱动程序值,每个数据库供应商都不同。这定义了在结果集中一次缓冲多少记录。例如:对于 Oracle,此值 10:通常太低,因为涉及从应用程序到数据库的大量读取操作;对于 Postgresql,这是无限的(整个结果集),然后太多了。您必须在速度和内存使用之间找到适当的平衡。
关于更新:
do {
YourType object = readFromKafkaQueue();
mybatisMapper.update(object);
} while (kafkaQueueHasMoreElements());
sqlSession.flushStatement(); // only when using ExecutorType.BATCH
最重要的是ExecutorType(这是SessionFactory.openSession()中的参数)或者ExecutorType.REUSE,这将允许只准备一次语句,而不是在每次迭代时使用默认的ExecutorType.SIMPLE或ExecutorType.BATCH,这将堆叠语句,实际上只在刷新时执行它们。
现在还需要考虑事务:这可能涉及提交 300 万次更新,或者可以分段。
【讨论】:
您需要为批处理创建单独的映射器实例。
这可能会有所帮助。
http://pretius.com/how-to-use-mybatis-effectively-perform-batch-db-operations/
【讨论】: