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

如何实现Java SpringBoot应用代理MQTT Broker的WebSocket 8083端口?

实现Spring Boot TCP代理转发MQTT WebSocket流量

要实现设备通过Spring Boot应用A的8083端口透明访问MQTT Broker的WebSocket服务,核心是做纯TCP层的双向转发代理——不需要解析WebSocket协议,只需要把设备发往A的所有TCP数据包原样转发给Broker,再把Broker的响应原路返回给设备。

下面是可行的实现方案和代码示例:

一、核心逻辑

代理的工作流程很简单:

  • 监听应用A的8083端口,接收设备的TCP连接请求
  • 与MQTT Broker的8083端口建立新的TCP连接
  • 双向绑定两个连接,实现数据的实时转发
  • 任意一端连接关闭时,同步关闭另一端连接

二、基于Netty的实现(推荐)

Netty是Java生态中成熟的NIO网络框架,适合构建高并发的代理服务,在Spring Boot中集成非常方便。

1. 添加Maven依赖

<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-all</artifactId>
    <version>4.1.100.Final</version>
</dependency>

2. 核心代理处理器

这个类负责处理单条连接的双向数据转发:

import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;

public class ProxyHandler extends ChannelInboundHandlerAdapter {
    private final Channel remoteChannel;

    public ProxyHandler(Channel remoteChannel) {
        this.remoteChannel = remoteChannel;
    }

    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) {
        if (remoteChannel.isActive()) {
            remoteChannel.writeAndFlush(msg);
        }
    }

    @Override
    public void channelInactive(ChannelHandlerContext ctx) {
        if (remoteChannel.isActive()) {
            remoteChannel.close();
        }
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
        cause.printStackTrace();
        ctx.close();
        remoteChannel.close();
    }
}

3. 客户端连接初始化器

用于初始化连接MQTT Broker的通道:

import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.socket.SocketChannel;

public class ClientChannelInitializer extends ChannelInitializer<SocketChannel> {
    private final Channel localChannel;

    public ClientChannelInitializer(Channel localChannel) {
        this.localChannel = localChannel;
    }

    @Override
    protected void initChannel(SocketChannel ch) {
        ChannelPipeline pipeline = ch.pipeline();
        pipeline.addLast(new ProxyHandler(localChannel));
    }
}

4. 代理服务器初始化器

负责启动本地监听端口,并为每个新连接创建到Broker的转发通道:

import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;

public class ProxyServerInitializer extends ChannelInitializer<SocketChannel> {
    private final String brokerHost;
    private final int brokerPort;

    public ProxyServerInitializer(String brokerHost, int brokerPort) {
        this.brokerHost = brokerHost;
        this.brokerPort = brokerPort;
    }

    @Override
    protected void initChannel(SocketChannel ch) {
        ChannelPipeline pipeline = ch.pipeline();
        pipeline.addLast(new ChannelInboundHandlerAdapter() {
            @Override
            public void channelActive(ChannelHandlerContext ctx) {
                EventLoopGroup group = new NioEventLoopGroup();
                Bootstrap bootstrap = new Bootstrap();
                bootstrap.group(group)
                        .channel(ctx.channel().getClass())
                        .handler(new ClientChannelInitializer(ctx.channel()));

                // 连接到MQTT Broker
                bootstrap.connect(brokerHost, brokerPort).addListener(future -> {
                    if (future.isSuccess()) {
                        Channel remoteChannel = ((Channel) future.getNow());
                        ctx.pipeline().addLast(new ProxyHandler(remoteChannel));
                    } else {
                        ctx.close();
                        group.shutdownGracefully();
                    }
                });
            }
        });
    }

    // 启动代理服务器
    public void start(int localPort) throws InterruptedException {
        EventLoopGroup bossGroup = new NioEventLoopGroup(1);
        EventLoopGroup workerGroup = new NioEventLoopGroup();
        try {
            io.netty.bootstrap.ServerBootstrap serverBootstrap = new io.netty.bootstrap.ServerBootstrap();
            serverBootstrap.group(bossGroup, workerGroup)
                    .channel(NioServerSocketChannel.class)
                    .childHandler(this);

            // 绑定本地端口并启动
            serverBootstrap.bind(localPort).sync().channel().closeFuture().sync();
        } finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }
    }
}

5. 在Spring Boot中启动代理

通过CommandLineRunner在应用启动时自动启动代理:

import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;

@Component
public class ProxyStarter implements CommandLineRunner {
    @Value("${mqtt.broker.host}")
    private String brokerHost;

    @Value("${mqtt.broker.port:8083}")
    private int brokerPort;

    @Value("${proxy.port:8083}")
    private int proxyPort;

    @Override
    public void run(String... args) throws InterruptedException {
        ProxyServerInitializer proxyServer = new ProxyServerInitializer(brokerHost, brokerPort);
        proxyServer.start(proxyPort);
    }
}

6. 配置参数

在application.properties中添加代理配置:

# MQTT Broker地址
mqtt.broker.host=your-broker-ip
# MQTT Broker的WebSocket端口
mqtt.broker.port=8083
# 代理监听的本地端口
proxy.port=8083

三、其他可选方案

  • Java NIO原生实现:不依赖框架,直接用ServerSocketChannel和SocketChannel实现双向转发,代码更底层,适合轻量低并发场景,但需要自己处理线程和资源管理。
  • 现成代理框架:比如使用Netty官方的netty-proxy模块,它提供了更完善的代理工具类,可以简化部分代码逻辑。

四、注意事项

  • 端口冲突:如果Spring Boot应用本身的服务端口也是8083,需要修改server.port调整应用端口,避免与代理端口冲突。
  • 资源泄漏:确保连接关闭时,EventLoopGroup等Netty资源能正确释放,防止内存泄漏。
  • 性能优化:可以根据并发量调整EventLoopGroup的线程数,Netty默认会根据CPU核心数分配线程,一般无需手动调整。

内容的提问来源于stack exchange,提问作者msl12

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 08:40:41