You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 23:00:55