我的mqtt协议和emqttd开源项目个人理解(9) - 集群和Mnesia源码分析

简介: 我的mqtt协议和emqttd开源项目个人理解(9) - 集群和Mnesia源码分析

学习mqtt协议和emqttd开源项目http://emqtt.com/

emqttd源码版本号是v1.1.3。http://emqtt.com/downloads/1113


一、先来看EMQ的文档定义:http://emqtt.com/docs/v1/cluster.html

emqttd集群设置管理¶


假设部署两台服务器s1.emqtt.io, s2.emqtt.io上部署集群:


节点名 主机名(FQDN) IP地址

emqttd@s1.emqtt.io 或emqttd@192.168.0.10 s1.emqtt.io 192.168.0.10

emqttd@s2.emqtt.io 或emqttd@192.168.0.20 s2.emqtt.io 192.168.0.20

Warning

节点名格式: Name@Host, Host必须是IP地址或FQDN(主机名.域名)


emqttd@s1.emqtt.io节点设置


emqttd/etc/vm.args:


-name emqttd@s1.emqtt.io



-name emqttd@192.168.0.10

Warning

节点启动加入集群后,节点名称不能变更。


emqttd@s2.emqtt.io节点设置


emqttd/etc/vm.args:


-name emqttd@s2.emqtt.io



-name emqttd@192.168.0.20

节点加入集群


启动两台节点后,emqttd@s2.emqtt.io上执行:


$ ./bin/emqttd_ctl cluster join emqttd@s1.emqtt.io


Join the cluster successfully.

Cluster status: [{running_nodes,['emqttd@s1.emqtt.io','emqttd@s2.emqtt.io']}]

或,emqttd@s1.emqtt.io上执行:


$ ./bin/emqttd_ctl cluster join emqttd@s2.emqtt.io


Join the cluster successfully.

Cluster status: [{running_nodes,['emqttd@s1.emqtt.io','emqttd@s2.emqtt.io']}]

任意节点上查询集群状态:


$ ./bin/emqttd_ctl cluster status


Cluster status: [{running_nodes,['emqttd@s1.emqtt.io','emqttd@s2.emqtt.io']}]

节点退出集群


节点退出集群,两种方式:


leave: 本节点退出集群

remove: 从集群删除其他节点

emqttd@s2.emqtt.io主动退出集群:


$ ./bin/emqttd_ctl cluster leave

或emqttd@s1.emqtt.io节点上,从集群删除emqttd@s2.emqtt.io节点:


$ ./bin/emqttd_ctl cluster remove emqttd@s2.emqtt.io

二、emqttd_ctl是怎么使用的?


-module(emqttd_cli).有定义要加载的命令


-export([status/1, broker/1, cluster/1, users/1, clients/1, sessions/1,
         routes/1, topics/1, subscriptions/1, plugins/1, bridges/1,
         listeners/1, vm/1, mnesia/1, trace/1]).
load() ->
    Cmds = [Fun || {Fun, _} <- ?MODULE:module_info(exports), is_cmd(Fun)],
    [emqttd_ctl:register_cmd(Cmd, {?MODULE, Cmd}, []) || Cmd <- Cmds].


三、加入集群后,子节点mnesia数据库怎么办?


mnesia数据库天然支持分布式集群。子节点加入之后就类似MySQL数据库主从备份一样,主节点和子节点mnesia会保持同步。来看源码:


-module(emqttd_mnesia).
%% @doc Join the mnesia cluster
-spec(join_cluster(node()) -> ok).
join_cluster(Node) when Node =/= node() ->
    %% Stop mnesia and delete schema first
    ensure_ok(ensure_stopped()),
    ensure_ok(delete_schema()),
    %% Start mnesia and cluster to node
    ensure_ok(ensure_started()),
    ensure_ok(connect(Node)),
    ensure_ok(copy_schema(node())),
    %% Copy tables
    copy_tables(),
    ensure_ok(wait_for(tables)).


子节点加入之后,会先删除自己的mnesia数据库和各个表,然后copy一份主节点的库,再copy各个表数据。

%% @doc Cluster with node.
-spec(connect(node()) -> ok | {error, any()}).
connect(Node) ->
    case mnesia:change_config(extra_db_nodes, [Node]) of
        {ok, [Node]} -> ok;
        {ok, []}     -> {error, {failed_to_connect_node, Node}};
        Error        -> Error
    end.
%% @doc Copy schema.
copy_schema(Node) ->
    case mnesia:change_table_copy_type(schema, Node, disc_copies) of
        {atomic, ok} ->
            ok;
        {aborted, {already_exists, schema, Node, disc_copies}} ->
            ok;
        {aborted, Error} ->
            {error, Error}
    end.
%% @doc Copy mnesia tables.
copy_tables() ->
    emqttd_boot:apply_module_attributes(copy_mnesia).


函数copy_tables(),会检索和执行emq工程目录下所有erl模块里面的mnesia(copy)函数,模块要求含有"-copy_mnesia({mnesia, [copy]})."关键字


例如:


-module(emqttd_backend).
mnesia(copy) ->
    ok = emqttd_mnesia:copy_table(retained_message),
    ok = emqttd_mnesia:copy_table(backend_subscription).
-module(emqttd_router).
-copy_mnesia({mnesia, [copy]}).
mnesia(copy) ->
    ok = emqttd_mnesia:copy_table(route, ram_copies).



EMQ工程目录下,有关键字-boot_mnesia({mnesia, [boot]}).和-copy_mnesia({mnesia, [copy]}).的模块是:


-module(emqttd_backend).

-module(emqttd_pubsub).

-module(emqttd_router).

-module(emqttd_server).

-module(emqttd_sm).

-module(emqttd_trie).


其中,emqttd_backend模块新建的数据库retained_message和backend_subscription是disc_copies类型的,其他模块是ram_copies类型的。




四、emqttd的mnesia初始化


1、-module(emqttd_app).

start(_StartType, _StartArgs) ->
    print_banner(),
    emqttd_mnesia:start(),


2、-module(emqttd_mnesia).

start() ->
    ensure_ok(ensure_data_dir()),
    ensure_ok(init_schema()),
    ok = mnesia:start(),
    init_tables(),
    wait_for(tables).
%% @doc Init mnesia schema or tables.
init_schema() ->
    case mnesia:system_info(extra_db_nodes) of
        []    -> mnesia:create_schema([node()]);
        [_|_] -> ok
    end.
%% @private
%% @doc Init mnesia tables.
init_tables() ->
    case mnesia:system_info(extra_db_nodes) of
        []    -> create_tables();
        [_|_] -> copy_tables()
    end.


3、数据库启动和拷贝的例子-module(emqttd_backend).

-boot_mnesia({mnesia, [boot]}).
-copy_mnesia({mnesia, [copy]}).
%% Mnesia callbacks
%%--------------------------------------------------------------------
mnesia(boot) ->
    ok = emqttd_mnesia:create_table(retained_message, [
                {type, ordered_set},
                {disc_copies, [node()]},
                {record_name, retained_message},
                {attributes, record_info(fields, retained_message)},
                {storage_properties, [{ets, [compressed]},
                                      {dets, [{auto_save, 1000}]}]}]),
    ok = emqttd_mnesia:create_table(backend_subscription, [
                {type, bag},
                {disc_copies, [node()]},
                {record_name, mqtt_subscription},
                {attributes, record_info(fields, mqtt_subscription)},
                {storage_properties, [{ets, [compressed]},
                                      {dets, [{auto_save, 5000}]}]}]);
mnesia(copy) ->
    ok = emqttd_mnesia:copy_table(retained_message),
    ok = emqttd_mnesia:copy_table(backend_subscription).


4、数据库启动和拷贝的例子-module(emqttd_router).


-boot_mnesia({mnesia, [boot]}).
-copy_mnesia({mnesia, [copy]}).
mnesia(boot) ->
    ok = emqttd_mnesia:create_table(route, [
                {type, bag},
                {ram_copies, [node()]},
                {record_name, mqtt_route},
                {attributes, record_info(fields, mqtt_route)}]);
mnesia(copy) ->
    ok = emqttd_mnesia:copy_table(route, ram_copies).

五、注意事项


1、如果EMQ所在服务器的IP地址是192.168.0.10


那么节点名称A:emqttd@192.168.0.10和节点名称B:emqttd@127.0.0.1是相同的意思,如果EMQ以节点A启动服务器,那么再以节点B启动是会失败的。


此时只能把A或B其中一个更名一下。即节点名格式: Name@Host里面的Name要加以区分。


2、集群的信息会记录在工程目录下,/rel/emqttd/data/mnesia/emqttd@172.16.6.161/schema.DAT


即,当子节点A连接了主节点B,集群信息会分别记录在schema.DAT。如果子节点A没有主动断开集群,下次重启时,仍然会主动连接主节点B。


★有几个遗留问题待确认,不知道EMQ V2版本有无修正:


问题(1)如果子节点A没有主动断开集群,下次重启时,如果B不存在,那么A就会启动失败!好可怕!


问题(2)A连接上B之后。A目录下的文件/rel/emqttd/data/mnesia/emqttd@172.16.6.161/retained_message.DCD和backend_subscription.DCD就自我删除了。以后也见不着了,彻底消失了,奇怪!请注意,这两个Mnesia数据库表类型是持久化,disc_copies。


★2018/05/17实测emq2.3.7,A主B从,结论如下:


(1)B join A之后,会主动删除B的Mnesia表,然后从A拷贝一份过来。B leave A之后,也会删除B的Mnesia表。


(2)B join A之后,A和B的Mnesia表会始终保持一致性。


添加或删除或更新A的表数据,B会同步。


添加或删除或更新B的表数据,A会同步。


(3)集群的信息会记录在工程目录下,/rel/emqttd/data/mnesia/emqttd@172.16.6.161/schema.DAT。A或B进程退出后,再次启动时,仍然保持集群状态。


3、常用命令


./emqttd console
./emqttd start
./emqttd stop
./emqttd_ctl cluster join emqttd@172.16.6.161
./emqttd_ctl cluster status
./emqttd_ctl cluster leave
werl -name firecat@127.0.0.1 -setcookie emqsecretcookie
observer:start().
./_rel/emqttd/bin/emqttd console
./_rel/emqttd/bin/emqttd start
./_rel/emqttd/bin/emqttd_ctl status
./_rel/emqttd/bin/emqttd stop
./_rel/emqttd/bin/emqttd_ctl cluster join emq@192.168.0.116
./_rel/emqttd/bin/emqttd_ctl cluster status
./_rel/emqttd/bin/emqttd_ctl cluster leave
相关实践学习
快速体验阿里云云消息队列RocketMQ版
本实验将带您快速体验使用云消息队列RocketMQ版Serverless系列实例进行获取接入点、创建Topic、创建订阅组、收发消息、查看消息轨迹和仪表盘。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
相关文章
|
11月前
|
数据采集 传感器 监控
Modbus 与 MQTT 协议兼容:MyEMS 的泛在能源数据采集技术实现
MyEMS深度融合Modbus与MQTT协议,破解能源数据采集中协议碎片化、网络异构、数据孤岛等难题。通过Modbus接入95%以上工业设备,实现现场数据精准“拉取”;依托MQTT构建高效物联网传输通道,支持多源数据主动“推送”与云端集成。边缘侧采集规整,中心侧汇聚分析,形成统一、可靠、低延迟的数据流。该架构兼具高兼容性、强扩展性与低运维成本,广泛应用于工业园区、商业楼宇及集团型企业,支撑实时监控、AI分析与跨系统融合,打造泛在互联的能源数据底座,助力企业实现全面智慧能源管理。
681 6
|
消息中间件 监控 RocketMQ
Docker部署RocketMQ5.2.0集群
本文详细介绍了如何使用Docker和Docker Compose部署RocketMQ 5.2.0集群。通过创建配置文件、启动集群和验证容器状态,您可以快速搭建起一个RocketMQ集群环境。希望本文能够帮助您更好地理解和应用RocketMQ,提高消息中间件的部署和管理效率。
2161 91
|
边缘计算 负载均衡 NoSQL
FreeMQTT Plus: 一个新型 MQTT Broker 集群的实现
FreeMQTT Plus 是一款基于 MQTT 协议的高性能消息中间件,采用分布式架构解决单点瓶颈问题。其核心由 Nginx 负载均衡器、黑(A)节点(MQTT Broker)、白(B)节点(消息路由)和日志(L)节点组成。通过无主从设计,支持高可用性、负载均衡与灵活扩展。针对会话同步、消息路由等挑战,FreeMQTT Plus 利用 MQTT5 特性定义元命令,实现节点间高效通信,无需依赖第三方组件。适用于物联网海量设备接入与高并发场景,为未来边缘计算和多级集群部署提供坚实基础。
2160 74
|
消息中间件 Apache 双11
Apache RocketMQ + “太乙” = 开源贡献新体验
Apache RocketMQ 是 Apache 顶级项目,源于阿里巴巴,历经多年双十一考验。RocketMQ 联合“太乙”平台启动开源竞赛,提供贡献价值评价与奖金激励(最高 5000 元),助力开发者成为社区核心成员。竞赛包含详尽教程与自动搭建环境,促进技术生态繁荣,推动分布式消息处理技术发展。欢迎加入,共创开源未来!
483 1
|
监控 安全 Java
Java 开发中基于 Spring Boot 3.2 框架集成 MQTT 5.0 协议实现消息推送与订阅功能的技术方案解析
本文介绍基于Spring Boot 3.2集成MQTT 5.0的消息推送与订阅技术方案,涵盖核心技术栈选型(Spring Boot、Eclipse Paho、HiveMQ)、项目搭建与配置、消息发布与订阅服务实现,以及在智能家居控制系统中的应用实例。同时,详细探讨了安全增强(TLS/SSL)、性能优化(异步处理与背压控制)、测试监控及生产环境部署方案,为构建高可用、高性能的消息通信系统提供全面指导。附资源下载链接:[https://pan.quark.cn/s/14fcf913bae6](https://pan.quark.cn/s/14fcf913bae6)。
2760 0
|
消息中间件 存储 Apache
恭喜 Apache RocketMQ 荣获 2024 开源创新榜单“年度开源项目”
恭喜 Apache RocketMQ 荣获 2024 开源创新榜单“年度开源项目”
404 1
|
消息中间件 数据管理 Serverless
阿里云消息队列 Apache RocketMQ 创新论文入选顶会 ACM FSE 2025
阿里云消息团队基于 Apache RocketMQ 构建 Serverless 消息系统,适配多种主流消息协议(如 RabbitMQ、MQTT 和 Kafka),成功解决了传统中间件在可伸缩性、成本及元数据管理等方面的难题,并据此实现 ApsaraMQ 全系列产品 Serverless 化,助力企业提效降本。
|
11月前
|
消息中间件 Java Kafka
消息队列比较:Spring 微服务中的 Kafka 与 RabbitMQ
本文深入解析了 Kafka 和 RabbitMQ 两大主流消息队列在 Spring 微服务中的应用与对比。内容涵盖消息队列的基本原理、Kafka 与 RabbitMQ 的核心概念、各自优势及典型用例,并结合 Spring 生态的集成方式,帮助开发者根据实际需求选择合适的消息中间件,提升系统解耦、可扩展性与可靠性。
714 1
消息队列比较:Spring 微服务中的 Kafka 与 RabbitMQ
|
消息中间件 JSON Java
开发者如何使用轻量消息队列MNS
【10月更文挑战第19天】开发者如何使用轻量消息队列MNS
1147 102