如何实现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
相关产品推荐
相关产品推荐

