如何使用Vert.x实现非阻塞批处理并支持集群分布式处理
Vert.x 实现分布式非阻塞批处理的可行方案
完全可以基于Vert.x实现你描述的批处理能力,以下是可落地的最优方案,可解决内存溢出、流协同、集群负载分摊三类核心问题:
核心实现逻辑
- 输入侧优先使用Vert.x原生带背压的流式API,避免全量加载数据:
- 文件读取用
vertx.fileSystem().open()返回的AsyncFile对象,它天然实现了Reactive Streams的Publisher接口,会根据下游处理速度自动控制读取速率,不会一次性加载全量文件到内存 - 数据库查询用Vert.x SQL Client的流式查询能力:调用
preparedQuery.createStream()方法即可按批次拉取结果集,拉取速度自动和下游处理能力对齐,不需要手动做分页逻辑
- 文件读取用
- 数据分发环节不要直接全量推送Event Bus,用点对点请求-响应模式做流控协同:
- 每从上游流拉取固定数量的记录(可根据业务配置,比如单次100条),向指定Event Bus地址发送带ack要求的消息
- 上游拉取逻辑绑定下游处理的响应回调:只有当前批次所有消息都收到处理完成的确认后,才触发上游拉取下一批数据,完全避免拉取速度快于处理速度导致的消息积压、内存溢出问题
集群负载分摊实现
- 把单条记录的处理逻辑封装成独立Verticle,在集群所有节点部署相同数量的实例,所有Verticle监听同一个Event Bus地址
- Vert.x集群模式内置负载均衡能力,会自动把Event Bus消息按轮询策略分发到不同节点的Verticle处理,不需要额外做服务发现、流量分发的开发
- 如果需要自定义负载规则,可实现
LoadBalancingStrategy接口,在集群启动时注册即可生效
可选优化
- 如果处理逻辑包含重CPU计算,可把处理Verticle配置为
WorkerVerticle,设置多实例部署,避免阻塞Event Loop - 对可靠性要求高的场景,可给Event Bus消息配置重试次数,搭配分布式Map做处理进度持久化,避免节点宕机丢失数据
内容的提问来源于stack exchange,提问作者Andrius
相关产品推荐
相关产品推荐

