首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >从 EMQX 到 DolphinDB 数据接入实操:百万 JSON 消息如何秒级入库

从 EMQX 到 DolphinDB 数据接入实操:百万 JSON 消息如何秒级入库

原创
作者头像
DolphinDB
发布2026-09-02 10:27:44
发布2026-09-02 10:27:44
1060
举报

100 万条工业数据,从 EMQX 到 DolphinDB 完成入库,需要多久?答案是:不到 1 秒。

在 8 月 20 日的直播中,我们与 EMQ 共同探讨了 FlowMQ + DolphinDB 构建轻量化实时数据底座的理念:让工业物联网数据不仅能够高效接入,更能够真正存下来、算起来。

在这套方案中,FlowMQ 负责数据的接入、分发与流转,DolphinDB 负责数据的存储、计算与分析,两者形成从设备数据接入到数据分析的完整链路。其中,EMQX 作为 FlowMQ 链路中的 MQTT Broker,承接设备通过 MQTT 协议上报的数据;DolphinDB 则通过 MQTT5 插件直接订阅 EMQX 消息,并在库内完成数据解析、计算与持久化。

这种协同方式将数据接入与后续存储、计算紧密衔接,减少额外中间组件带来的数据搬运和链路延迟,让高频、多源的工业数据能够更快转化为可分析、可持续积累的数据资产,为工业物联网提供一套轻量、高效的实时数据底座。

下面我们将代码详细拆解,从 MQTT 消息订阅、JSON 数据解析,到分布式入库与查询验证,逐步展示 DolphinDB 如何与 EMQX 协同,构建一条极简、高效、可扩展的工业数据接入与分析链路。

极简链路,一步到位

在典型的工业物联网场景中,设备数据通过 MQTT 协议上报至 EMQX Broker。传统方案往往需要引入 Kafka、Flink 等多套组件,来完成数据的解析、清洗、计算与持久化写入,链路冗长、组件多、运维复杂,同时,跨组件的数据复制与格式转换也会带来额外延迟,影响工业场景对实时性的要求

工业数据真正需要的,并不是更多的中间环节,而是从数据接入到计算、存储的更短路径。 基于这一思路,DolphinDB 与 EMQX 推出联合解决方案,将数据链路进一步压缩为:设备 → EMQX → DolphinDB。

DolphinDB 直接订阅 EMQX 中的 MQTT 消息,在库内完成数据解析、流式处理与分布式持久化,减少中间环节,让数据更快进入可分析状态。

  • 数据采集:NeuronEX 等采集软件将工业设备数据发布至 EMQX;
  • 实时订阅:DolphinDB 通过 MQTT5 插件 直接向 EMQX 发起订阅,实时拉取消息;
  • 流式写入:数据经由流表实时接入,通过内置的 parseJsonTable 函数完成 JSON 解析与结构化映射;
  • 持久化存储:解析后的数据批量写入 DolphinDB 分布式数据库,为后续高性能时序分析奠定底座。

整个链路中,DolphinDB 承担了消息订阅、三大任务,无需额外引入消息中间件或 ETL 工具,极大降低了系统复杂度。

实现步骤与关键逻辑

以下是接入方案的核心步骤。

1、安装与加载 MQTT5 Plugin

DolphinDB 为 MQTT 协议提供了原生插件支持,要求 Server 版本 2.00.10 及以上:

代码语言:javascript
复制
installPlugin("mqtt5")
loadPlugin("mqtt5")

2、创建流式消息接收表

在 DolphinDB 中创建一个流表,用于暂存从 EMQX 订阅到的原始消息:

代码语言:javascript
复制
share streamTable(1000000:0, `topic`message, [STRING, STRING]) as mqtt5Raw

3、向 EMQX 发起订阅

配置订阅参数,指定回调函数将消息写入流表:

代码语言:javascript
复制
config = dict(STRING, ANY)
config['qos'] = 0
config['clientID'] = "ddb_mqtt5_benchmark"
config['mqttVersion'] = 4
sub5 = mqtt5::subscribe(
    "tcp://183.134.101.140:1883",
    "benchmark/#",
    onMessage{mqtt5Raw},
    config
)
mqtt5::getSubscriberStat()

执⾏过后可以从 DolphinDB 侧看到连接信息:

同时也可以在 EMQX Dashboard 上⾯看到连接信息,前者为与 NeuronEX 的连接,后者为与 DolphinDB 的连接:

4、创建分布式数据库表

建立分区表以支撑海量数据存储:

代码语言:javascript
复制
dbPath = "dfs://IndustrialDB"
tableName = `benchmark_industrial
topicPartitions = string(
    "benchmark/" + string(0..15)
)
db = database(
    directory=dbPath,
    partitionType=VALUE,
    partitionScheme=topicPartitions,
    engine="OLAP"
)

//dropDatabase(dbPath)
schemaTable = table(
    1:0,
    colNames,
    colTypes
)
industrialTable = createPartitionedTable(
    dbHandle=db,
    table=schemaTable,
    tableName=tableName,
    partitionColumns=`topic
)

5、流表数据自动解析与入库

通过 subscribeTable 订阅流表,自定义解析函数将 JSON 消息映射为结构化数据,并支持批量写入以提升吞吐:

代码语言:javascript
复制
setStreamTableFilterColumn(
    mqtt5Raw,
    `topic
)
jsonSchema = table(
    ["topic","run_id","seq","event_ts","device_id","point_id","metric","value","unit","quality"] as name,
    ["STRING", "STRING","LONG","LONG","STRING","STRING","STRING","DOUBLE","STRING","INT"] as type
)
def parser(mutable outputTable,schema,msg)
{
    topicValue  = msg.topic[0]
    jsonText = concat(msg.message, "\n")           // 多条 MQTT message 拼成一段连续的 NDJSON
    parsedData = parseJsonTable(jsonText, schema)  // 一次解析成多行表
    n = size(parsedData)
    update!(parsedData,`topic,take(topicValue,n))
    append!(outputTable, parsedData)               // 一次批量入库

}

for(i in 0..(parallel - 1))
{
    topicName = "benchmark/" + string(i)
    subscribeTable(
        tableName="mqtt5Raw",
        actionName="parseMqttJson_" + string(i),
        offset=0,
        handler=parser{
            benchmark_industrial,
            jsonSchema
        },
        msgAsTable=true,
        batchSize=1000,
        throttle=0.1,
        hash=i,
        filter=[topicName]
    )
}

分区设计要点,如何跑出最快速度

为进一步发挥链路下游的吞吐能力,我们围绕多 worker 并行消费 + 分区绑定进行优化设计,通过主题分发、分区映射与写入隔离,提升从 EMQX 到 DolphinDB 的整体入库效率。

具体设计包括以下三点:

  • 按业务流量将消息拆分至 benchmark/0benchmark/15 共 16 个主题,形成可并行消费的数据通道
  • 分布式表按 seq 列进行哈希分区,使数据能够稳定落入对应分区
  • 启动 16 个 worker,每个 worker 通过 hash=i 绑定到特定主题进行订阅消费,实现主题、worker 与分区的协同映射

优化后的核心效果在于:每个 worker 处理的 seq 取值范围互不重叠,形成 worker 与分区的一一映射。这样既能充分利用多 worker 并行消费能力,又能让每个 worker 独占对应分区写入,避免分区竞争,从而显著提升 DolphinDB 分布式并行写入效率。

核心实现代码如下:

代码语言:javascript
复制
// 设置流表按 topic 列进行过滤分发
setStreamTableFilterColumn(
    mqtt5Raw,
    `topic
)

// 启动 16 个 worker,每个订阅一个主题
// 通过 hash=i 将 worker 与分区绑定,实现写入隔离,避免分区竞争
for(i in 0..15) {
    topicName = "benchmark/" + string(i)
    subscribeTable(
        tableName="mqtt5Raw",
        actionName="parseMqttJson_" + string(i),
        offset=0,
        handler=parser{benchmark_industrial, jsonSchema},
        msgAsTable=true,
        batchSize=1000,    // 多 worker 并行时建议调小批次,避免内存压力
        throttle=0.1,
        hash=i,            // 关键:将订阅绑定到指定分区
        filter=[topicName]
    )
}

核心价值:不止于接入,更在于融合

这套方案之所以能够成为直播观众关注的焦点,源于以下三个层面的核心价值:

① 原生集成,简化接入链路

借助 MQTT5 插件,DolphinDB 可以直接订阅 EMQX 的消息,并在库内完成 JSON 解析、流式处理与分布式入库。相比再额外叠加 Kafka、Flink 或自研消费客户端的方案,这一路径将设备数据从接入到落库的链路进一步缩短,减少跨组件转发与格式转换带来的数据搬运。

② 流批一体,实时入库与分析兼得

数据经由流表进入系统后,一方面通过订阅实现持续落盘,另一方面可同时挂载流计算引擎进行实时告警、异常检测等处理。一份数据,既满足存储需求,也支撑实时分析。

③ 性能实测:百万数据秒级入库

为验证 FlowMQ + DolphinDB 完整链路的实际性能,本次对 EMQX 与 DolphinDB 进行了多次联合测试,测试环境如下:

  • EMQX:v6.2.2 试用版,单节点部署
  • DolphinDB:三节点高可用集群部署,通过 MQTT5 插件直接订阅 EMQX 消息并完成入库
  • 服务器配置:Intel® Xeon® Silver 4214 @ 2.20 GHz,x86 架构,单机 24 个物理核心、48 个逻辑线程
  • 测试数据:1000 万行

基于 EMQX 与 DolphinDB 的多次联合测试,

测试环节

峰值性能

EMQX → DolphinDB 流表

450 万行/秒

EMQX → 流表 → 分布式表(完整链路)

120 万行/秒

需要注意的是,完整链路的吞吐瓶颈主要在于流表写入分布式表的环节,但 120 万行/秒的性能已远超真实工业场景中 NeuronEX 等采集软件 10 万行/秒的产出上限,具备充足的性能冗余。

写在最后

在直播中我们反复强调一个理念:AI 竞争的本质是数据实时性的竞争。而实时数据竞争的第一步,就是让数据能够以最低的延迟、最简单的路径,从设备流向分析引擎。

从 EMQX 到 DolphinDB,这套方案用最少的组件、最简洁的代码,打通了工业数据从采集到存储分析的最后一公里——120 万行/秒的吞吐能力,足以支撑最苛刻的工业实时采集场景。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 极简链路,一步到位
  • 实现步骤与关键逻辑
    • 1、安装与加载 MQTT5 Plugin
    • 2、创建流式消息接收表
    • 3、向 EMQX 发起订阅
    • 4、创建分布式数据库表
    • 5、流表数据自动解析与入库
  • 分区设计要点,如何跑出最快速度
  • 核心价值:不止于接入,更在于融合
  • 写在最后
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档