基于Apache Camel Netty实现带握手机制与集群轮询的TCP Publisher和Subscriber的方案咨询
基于Apache Camel Netty实现带握手机制与集群轮询的TCP Publisher和Subscriber的方案咨询
嘿,这个场景我之前刚好基于Apache Camel Netty落地过,刚好能给你梳理下具体的实现思路,分Subscriber和Publisher两部分来说,应该能解决你的问题:
一、TCP Subscriber 实现(握手重试+集群轮询)
你的核心需求是握手失败重试、失败自动切到下一个服务器,Camel的负载均衡和流程控制组件刚好能搞定这个:
1. 集群轮询的连接管理
首先把所有目标服务器地址配置成逗号分隔的列表(比如在application.properties里写tcp.servers=server1:8080,server2:8080,server3:8080),然后用Camel的loadBalance().roundRobin()做轮询负载均衡,配合失败转移策略,一旦当前服务器握手失败/连接异常,自动切换到下一个节点。
2. 握手机制的流程控制
用Camel的choice()分支判断握手响应,结合retry()或者循环逻辑处理NOK的重试场景,成功后进入数据接收流程。这里给你个Java DSL的路由示例:
// 入口路由:负责集群轮询和失败转移 from("netty:tcp://{{tcp.servers}}?sync=true&clientMode=true&failOnException=true") .loadBalance().roundRobin().failOver(3, true, true) // 最多失败3次切换,支持异常转移 .to("direct:handshake-handler") .end(); // 握手处理路由 from("direct:handshake-handler") // 发送初始化命令 .setBody(constant("YOUR_INITIAL_HANDSHAKE_CMD")) .to("netty:tcp://${header.CamelNettyRemoteAddress}?sync=true") // 复用当前连接的服务器地址 .choice() // 握手成功,进入数据接收 .when(body().isEqualTo("OK")) .log("Handshake succeeded, start receiving data") .to("direct:data-receiver") // 握手失败,重试当前流程 .when(body().isEqualTo("NOK")) .log("Handshake failed with NOK, retrying...") .retry(3) // 最多重试3次,超过则触发失败转移到下一个服务器 .to("direct:handshake-handler") // 未知响应,直接触发失败转移 .otherwise() .log("Unknown response: ${body}, failover to next server") .throwException(new RuntimeException("Invalid handshake response")) .endChoice(); // 数据接收路由 from("direct:data-receiver") .log("Received data from server: ${body}") // 这里写你的业务数据处理逻辑 .process(exchange -> { String receivedData = exchange.getIn().getBody(String.class); // 处理数据... });
关键注意点
- 用
sync=true开启同步请求响应,因为握手需要等待服务器的明确回复 failOnException=true确保异常能触发负载均衡的失败转移- 重试次数别设成无限循环,避免死耗在一个故障节点上
二、TCP Publisher 实现(单连接限制+握手校验)
你的需求是只允许一个客户端连接、校验命令后发送数据,Netty的参数和Camel的处理器可以轻松实现:
1. 单连接限制
在Netty端点配置里加maxConnections=1,这样Camel Netty会自动拒绝超过1个的连接请求,或者你也可以配合allowDefaultCodec=false自定义连接拦截逻辑。
2. 握手校验与数据发送流程
当客户端发送命令过来,先在处理器里校验有效性,无效就关闭连接;有效则返回OK,然后启动循环发送数据的流程。Java DSL示例:
import org.apache.camel.component.netty.NettyConstants; // 监听客户端连接的入口路由 from("netty:tcp://0.0.0.0:{{tcp.pub.port}}?maxConnections=1&sync=true&allowDefaultCodec=true") .process(exchange -> { String clientCmd = exchange.getIn().getBody(String.class); // 自定义命令校验逻辑 if (!isValidClientCommand(clientCmd)) { // 标记关闭当前连接 exchange.getIn().setHeader(NettyConstants.NETTY_CLOSE_CHANNEL, true); exchange.getIn().setBody("NOK"); log.warn("Invalid client command, closing connection"); } else { exchange.getIn().setBody("OK"); // 触发后台数据发送流程,用单独的线程避免阻塞握手响应 exchange.getContext().createProducerTemplate().asyncSend("direct:data-sender", exchange); } }); // 后台数据发送路由 from("direct:data-sender") .loopDoWhile(constant(true)) // 循环发送直到连接断开 .delay(1000) // 模拟数据发送间隔,根据你的业务调整 .setBody(method(YourDataGenerator.class, "getNextData")) // 自定义数据生成逻辑 .to("netty:tcp://${header.CamelNettyRemoteAddress}?sync=false") .onException(IOException.class) .log("Client disconnected, stopping data sending") .stop() // 连接断开时停止循环 .end(); // 自定义校验方法 private boolean isValidClientCommand(String cmd) { // 这里写你的命令合法性校验逻辑,比如检查格式、签名等 return cmd != null && cmd.startsWith("VALID_PREFIX_"); }
关键注意点
- 用
asyncSend()异步启动数据发送,避免阻塞握手响应的回复 NettyConstants.NETTY_CLOSE_CHANNEL头是Camel Netty的内置标识,设置为true会自动关闭当前通道- 数据发送路由里要捕获连接异常,及时停止循环,避免无效的发送尝试
三、通用优化建议
- 连接状态监控:可以给Subscriber和Publisher加Camel的健康检查组件,监控连接状态和握手成功率
- 超时配置:在Netty端点加
connectTimeout=3000(3秒)、requestTimeout=5000等参数,避免长期阻塞 - 日志埋点:在握手、失败转移、连接断开等关键节点加日志,方便排查问题
内容来源于stack exchange
相关产品推荐
相关产品推荐

