Spring Integration:如何手动控制客户端模式下TcpReceivingChannelAdapter的连接与断开?
Spring Boot Kotlin TCP客户端手动连接控制方案
需求说明
- 作为TCP客户端连接远程服务器
- 远程服务器主动发送请求,应用按需回复(telnet连接时,未回复前会收到服务器初始消息)
- 并非所有请求都需回复,Spring Gateway方案不适用
- 需通过编程方式(如REST端点)手动控制连接的建立与断开,启动时不自动连接
现有代码问题
目前简化代码功能正常,但存在以下问题:
- Spring Boot启动时会自动建立TCP连接,若远程服务器未监听则抛出
connection refused异常 - 无法通过编程方式手动触发连接或断开连接
- 尝试设置
isAutoStartup=false未生效,推测与连接轮询机制有关
现有简化代码:
@Bean fun inboundFlow(connectionFactory: AbstractClientConnectionFactory): IntegrationFlow { val inbound = TcpReceivingChannelAdapter() inbound.isClientMode = true inbound.isAutoStartup = false inbound.setConnectionFactory(connectionFactory) return IntegrationFlow.from(inbound) .channel(requests()) .get() } @Bean fun outboundFlow(connectionFactory: AbstractClientConnectionFactory): IntegrationFlow { val outbound = TcpSendingMessageHandler() outbound.isClientMode = true outbound.setConnectionFactory(connectionFactory) return IntegrationFlow.from(replies()) .transform(MyTransformer::toByteArray) .handle(outbound) .get() } @Bean fun connectionFactory(): AbstractClientConnectionFactory = Tcp.nioClient(hostname, port) .singleUseConnections(false) .get() @Bean fun myFlow(connectionFactory: AbstractClientConnectionFactory, myReceiverService: MyReceiverService): IntegrationFlow? { return IntegrationFlow.from(requests()) .transform(MessageTransformer::transformMyMessage) .handle { msg: Message<MyMessage> -> myReceiverService.handleMessage(msg) } .get() }
尝试过的无效方案
- 设置
isAutoStartup=false:未生效,连接仍会自动建立 - Spring Profiles:需要重启应用,会重置其他服务状态,不可接受
- 控制通道:不清楚如何应用到当前场景
最终解决方案
通过IntegrationFlowContext动态注册和移除与连接关联的IntegrationFlow,实现手动控制连接的建立与断开。核心思路是在需要连接时才创建并注册相关的Flow,断开时移除这些Flow。
解决方案代码:
class MyCustomIntegrationFlow( private val flowContext: IntegrationFlowContext, @Value("${tcp.client.hostname}") private val hostname: String, @Value("${tcp.client.port}") private val port: Int, private val requests: MessageChannel, private val replies: MessageChannel ) { private val logger = KotlinLogging.logger {} private val inboundFlowId = "inboundFlowId" private val outBoundFlowId = "outBoundFlow" private val myCustomFlowId = "myCustomFlowId" private var started = false fun connect() { if (!started) { started = true val connectionFactory = connectionFactory() val inboundFlow = inboundFlow(connectionFactory) val outBoundFlow = outboundFlow(connectionFactory) val myCustomFlow = myCustomFlow() logger.info { "Registering inbound flow" } flowContext.registration(inboundFlow) .id(inboundFlowId) .register() logger.info { "Registering outbound flow" } flowContext.registration(outBoundFlow) .id(outBoundFlowId) .register() logger.info { "Registering Custom flow" } flowContext.registration(myCustomFlow) .id(myCustomFlowId) .register() } } fun disconnect() { if (started) { logger.info { "Removing Custom flow" } flowContext.remove(myCustomFlowId) logger.info { "Removing inbound flow" } flowContext.remove(inboundFlowId) logger.info { "Removing outbound flow" } flowContext.remove(outBoundFlowId) started = false } } private fun myCustomFlow(): IntegrationFlow { return IntegrationFlow.from(requests) .transform( /* 自定义转换逻辑 */) .handle {/* 自定义消息处理逻辑 */} .get() } private fun inboundFlow(connectionFactory: TcpNioClientConnectionFactorySpec): IntegrationFlow { val inbound = Tcp.inboundAdapter(connectionFactory) .clientMode(true) return IntegrationFlow.from(inbound) .channel(requests) .get() } // 重要:保持工厂为Spec类型,确保Spring连接事件正常工作 private fun connectionFactory(): TcpNioClientConnectionFactorySpec = Tcp.nioClient(hostname, port) .singleUseConnections(false) .serializer(TcpCodecs.lengthHeader(4)) .deserializer(TcpCodecs.lengthHeader(4)) private fun outboundFlow(connectionFactory: TcpNioClientConnectionFactorySpec): IntegrationFlow { val outbound = Tcp.outboundAdapter(connectionFactory) .clientMode(true) return IntegrationFlow.from(replies) .transform(/* 自定义转换调用 */) .handle(outbound) .get() } }
该方案成功实现了:
- Spring Boot启动时不自动建立TCP连接
- 通过调用
connect()和disconnect()方法(可封装为REST端点)手动控制连接的建立与断开
内容的提问来源于stack exchange,提问作者DuckTyper
相关产品推荐
相关产品推荐

