基于Kafka的Oracle到MongoDB数据迁移线程异常致丢数问题
线程分配逻辑问题导致Oracle到MongoDB迁移缺失100万条记录的排查与解决
嘿,我之前处理过类似的批量数据拆分迁移的坑,咱们从你给出的代码片段和场景入手,一步步揪出问题所在。
核心问题推测:线程范围计算错误(尤其是Oracle rownum的特性)
你提到分10个线程处理1000万条数据,但总有一个线程没运行,本质上大概率是数据拆分的范围计算逻辑有漏洞,结合Oracle的特性,最可能的点有这几个:
1. Oracle rownum的起始值误区
Oracle的rownum是从1开始计数的伪列,但你的代码里初始minRownum = 0,这会直接导致查询逻辑出错:
- 如果第一个线程的查询条件是
rownum between 0 and 1000000,Oracle会自动忽略rownum < 1的条件,实际只查了1-1000000的100万条(这部分是对的) - 但后续线程的范围如果基于0累加,比如第二个线程是
1000000-1999999,Oracle里rownum是查询结果的行号,这个范围会匹配不到任何数据(因为rownum不会大于当前查询结果的行数),最后一个线程的范围可能完全落在无效区间,直接导致线程无数据可处理,看起来像是“未运行”。
2. 循环边界的计数错误
你的代码里for(int i=minRownum;i<...的终止条件如果写得不对,会直接少生成一个线程。比如如果循环条件是i < totalRec/recInThread,但实际应该循环10次,却只循环了9次,自然少一个线程。
修正后的代码示例
我给你调整了线程拆分的逻辑,兼顾Oracle的特性和健壮性:
int totalRec = countNoOfRecordsToBeProcessed; int threadCount = 10; int recPerThread = totalRec / threadCount; // 处理非10整数倍的情况(虽然你这里是1000万,但保留逻辑更通用) int remainder = totalRec % threadCount; System.out.println("oracle " + new Date()); for (int i = 0; i < threadCount; i++) { // 从1开始计算rownum范围,适配Oracle特性 int minRownum = i * recPerThread + 1; int maxRownum = (i + 1) * recPerThread; // 最后一个线程把余数补上(如果有),确保所有数据都被覆盖 if (i == threadCount - 1) { maxRownum += remainder; } // 打印每个线程的范围,方便排查 System.out.printf("线程%d处理范围:%d - %d%n", i+1, minRownum, maxRownum); // 创建并启动线程,传入正确的查询范围 new Thread(() -> { // 执行Oracle查询:SELECT * FROM your_table WHERE rownum BETWEEN ? AND ? // 注意:如果是复杂查询,建议用子查询先排序再取rownum,避免数据顺序混乱 // 然后发送到Kafka,再同步到MongoDB }).start(); }
额外排查步骤
如果调整代码后还有问题,可以做这几件事:
- 打印每个线程的
minRownum和maxRownum,确认10个线程的范围完全覆盖1-10000000 - 手动执行最后一个线程的SQL语句,验证是否能查询到100万条记录
- 给线程添加日志,查看未运行的线程是否抛出了SQL异常、连接异常等
- 检查Kafka的生产监控,确认每个线程都在正常发送消息;再查MongoDB的写入日志,看是否有线程的写入被中断
内容的提问来源于stack exchange,提问作者user1708054
相关产品推荐
相关产品推荐

