当前位置:首页 > 物联网 > 智能应用
[导读]十万级设备后台别用Tomcat开HTTP长轮询,正确结构是“轻量接入网关收MQTT/CoAP/Modbus-TCP→协议归一→异步认证→消息入队列→业务服务做规则、存储、命令下发”。Java侧用Netty做接入、Spring Boot做治理,既轻又易维护;纯MQTT小系统也可用Moquette内嵌Broker,大并发再切自建Netty网关或EMQX集群。


十万级设备后台别用Tomcat开HTTP长轮询,正确结构是“轻量接入网关收MQTT/CoAP/Modbus-TCP→协议归一→异步认证→消息入队列→业务服务做规则、存储、命令下发”。Java侧用Netty做接入、Spring Boot做治理,既轻又易维护;纯MQTT小系统也可用Moquette内嵌Broker,大并发再切自建Netty网关或EMQX集群。

一、Netty接入网关骨架

MQTT长连接用主从Reactor,boss只accept,worker跑IO;协议解析走netty-codec-mqtt,心跳、限流、鉴权依次入Pipeline,避免IO线程做DB。

public class IotMqttServer {

   public void start(int port) throws Exception {

       EventLoopGroup boss = new NioEventLoopGroup(1);

       EventLoopGroup worker = new NioEventLoopGroup(

           Math.min(16, Runtime.getRuntime().availableProcessors() * 2));

       try {

           ServerBootstrap b = new ServerBootstrap();

           b.group(boss, worker)

            .channel(NioServerSocketChannel.class)

            .option(ChannelOption.SO_BACKLOG, 2048)

            .childOption(ChannelOption.TCP_NODELAY, true)

            .childOption(ChannelOption.SO_KEEPALIVE, true)

            .childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)

            .childHandler(new ChannelInitializer<SocketChannel>() {

                protected void initChannel(SocketChannel ch) {

                    ChannelPipeline p = ch.pipeline();

                    p.addLast("decoder", new MqttDecoder());

                    p.addLast("encoder", new MqttEncoder());

                    p.addLast(new IdleStateHandler(0, 0, 120));

                    p.addLast(new QuotaHandler());      // IP/租户连接数限流

                    p.addLast(new AsyncAuthHandler());  // 异步查Redis/认证服务

                    p.addLast(new PublishDispatchHandler()); // 转业务线程池

                }

            });

           b.bind(port).sync().channel().closeFuture().sync();

       } finally { boss.shutdownGracefully(); worker.shutdownGracefully(); }

   }

}

IdleStateHandler按业务设读写空闲,超过即断链并清在线状态;PooledByteBuf减少十万连接下的缓冲分配压力。

二、异步认证与会话生命周期

设备CONNECT带productKey、deviceKey、token,Handler只做轻校验,重校验提交业务Executor,不阻塞EventLoop:

public class AsyncAuthHandler extends SimpleChannelInboundHandler<MqttMessage> {

   private final ExecutorService authPool = Executors.newFixedThreadPool(8);

   private final DeviceRegistry registry;


   protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) {

       if (!(msg instanceof MqttConnectMessage m)) { ctx.fireChannelRead(msg); return; }

       String cid = m.payload().clientIdentifier();

       String user = m.payload().userName();

       String pwd = new String(m.payload().passwordInBytes(), StandardCharsets.UTF_8);

       authPool.submit(() -> {

           Device d = registry.authenticate(user, cid, pwd);

           if (d == null) { ctx.close(); return; }

           registry.markOnline(d, ctx.channel());

           ctx.writeAndFlush(buildConnAck());

       });

   }

}

在线状态存Redis Hash/Set,key按tenant:product:device,带心跳续期;断线、idle超时、全连接关闭都触发markOffline,命令服务查不在线改离线队列。会话用cleanSession=false可保留订阅与QoS1/2离线消息,但十万设备要限会话过期时间,避免Broker内存膨胀。

三、发布分流与业务解耦

PUBLISH不在IO线程写库,封装成IngestEvent投到业务线程池或Kafka/RabbitMQ:

public class PublishDispatchHandler extends SimpleChannelInboundHandler<MqttMessage> {

   private final ExecutorService biz = new DefaultEventExecutorGroup(16);

   protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) {

       if (!(msg instanceof MqttPublishMessage p)) { ctx.fireChannelRead(msg); return; }

       String topic = p.variableHeader().topicName();

       byte[] payload = ByteBufUtil.getBytes(p.payload());

       biz.execute(() -> {

           IngestEvent e = TopicRouter.parse(topic, payload); // up/{tenant}/{pk}/{dev}/{type}

           RuleEngine.run(e);

           TsStore.asyncWrite(e);

           CommandTracker.ackIfMatch(e);

       });

       if (p.fixedHeader().qos() == MqttQoS.AT_LEAST_ONCE)

           ctx.writeAndFlush(buildPubAck(p.variableHeader().packetId()));

   }

}

Topic按租户、产品、设备、消息类型分层,路由、权限、统计都跟着走;高频遥测QoS0、控制回执QoS1、关键配置可用QoS2并在业务层做幂等。

四、Spring Boot轻量治理

网关只管连接,Spring Boot管设备注册、命令、证书、监控:

@Component

public class CommandSender {

   public void send(String deviceId, String json) {

       String topic = "down/" + registry.topicOf(deviceId) + "/cmd";

       if (registry.isOnline(deviceId))

           mqttGateway.publish(topic, json, 1);

       else

           offlineCmdStore.enqueue(deviceId, topic, json);

   }

}

用Spring Integration或Paho outbound连外部MQTT Broker均可;配置、指标、Nacos/Apollo、Prometheus统一接Spring Boot,Netty负责高性能网络,Spring负责工程化。

五、十万级容量落地

单节点不是硬扛十万长连,而是按“接入分片+状态外置+数据异步”横向扩展:

Linux调fd与队列:ulimit -n 1000000,net.core.somaxconn、tcp_keepalive、SO_BACKLOG配合;

JVM用G1/ZGC,固定Xms=Xmx,MaxDirectMemory给池化ByteBuf留余量,开Netty泄漏检测SIMPLE;

设备状态、订阅、限流计数放Redis集群,接入节点无状态,扩节点只加Netty网关;

按tenant或productKey做一致性哈希分片,命令下发根据设备所在节点路由,避免全广播;

弱网设PINGREQ/PINGRESP、指数重连,离线命令带messageId与超时重投,设备回reply主题关单;

安全用每设备独立凭证,传输TLS/mTLS,平台侧JWT短期token,Broker做ACL限制主题范围。

按“Netty接入+异步鉴权+Redis会话+MQ异步处理+分片路由”的轻量组合,十台接入节点管十万设备更现实:单节点故障只影响分片,扩容加网关节例,业务层用Spring Boot做设备、规则、命令和可观测,整体比单体HTTP后台更适合嵌入式物联网长期运行。



本站声明: 本文章由作者或相关机构授权发布,目的在于传递更多信息,并不代表本站赞同其观点,本站亦不保证或承诺内容真实性等。需要转载请联系该专栏作者,如若文章内容侵犯您的权益,请及时联系本站删除( 邮箱:macysun@21ic.com )。
关闭