MySQL流处理执行5次后退出问题排查求助
MySQL Stream Processing Stops After 5 Iterations (With Unprocessed Data Left)
嘿,我帮你揪出问题所在了!你的代码里有几个关键错误,刚好解释了为什么只跑5次就停了——咱们一步步来梳理:
问题根源
- 错误分支没调用
callback():看你的代码,当Google Maps请求出错时,你只打了日志、恢复了连接,但没调用callback()。Transform流需要这个回调来标记「当前数据处理完了,可以收下一个」,没调用的话,流会一直认为当前任务没结束,缓冲区占满后就停止接收新数据了。 - 画蛇添足的
pause()/resume():你手动暂停/恢复MySQL连接的操作其实干扰了流的背压机制。Transform流本身会自动处理:当你在transform里没调用callback时,它会暂停读取上游的MySQL数据,直到你调用回调才继续。手动操作反而打乱了这个逻辑。 - MySQL流默认的
highWaterMark=5:这就是为什么刚好是5次!MySQL的查询流默认缓冲区大小是5个对象,当你没释放缓冲区(没调用callback),缓冲区一满,流就停止读取后续数据了。
修正后的代码
function AddDataSync() { var query = connection.query('select officename, divisionname from table') .stream() .pipe(stream.Transform({ objectMode: true, transform: function(data, encoding, callback) { // 移除不必要的手动暂停/恢复,交给流的背压机制处理 googleMapsClient.geocode({ address: data.officename + ' ' + data.divisionname }, function(err, response) { if (!err) { console.log(response.json.results[0].geometry.location); } else { console.log(err); } // 无论成功还是失败,必须调用callback! // 这是告诉Transform流:当前数据处理完毕,可以继续了 callback(); }); } })).on('finish',function() { console.log('done'); }); }
修正逻辑说明
- 去掉
connection.pause()和resume():让Transform流的背压机制自动控制读取速度,避免内存过载,也不会出现提前停止的情况。 - 所有分支都调用
callback():确保每一条数据处理完(不管成功失败),流都会继续读取下一条,直到数据库里的所有数据都处理完毕,最后触发finish事件打印「done」。
如果之后你想控制并发请求数(比如避免一下子发太多Google Maps请求被限流),可以再用队列或者async库做进一步优化,但当前问题的核心就是这两点啦!
内容的提问来源于stack exchange,提问作者Raghavendar Reddy
相关产品推荐
相关产品推荐

