后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载本篇技术指南围绕 emqx_bridge_tablestore 模块展开介绍 EMQX 如何通过连接器Connector 动作Action双层模型与阿里云表格存储Tablestore即 OTS集成将规则引擎筛选出的 MQTT 消息批量写入 Tablestore 时序表。读完本文你将掌握 Tablestore 桥接的完整配置参数、占位符模板渲染机制、时间戳与元数据更新语义、健康检查原理以及桥接的创建、管理与验证方法。一、桥接概览EMQX 与阿里云表格存储的数据通道emqx_bridge_tablestore是 EMQX 企业版数据集成体系中的一员用于将 EMQX 与阿里云表格存储Tablestore / OTS阿里云提供的 NoSQL 多模型数据库支持宽表、时序等数据模型打通。其核心职责是接收来自 EMQX 规则引擎的消息将其构造成时序数据行并写入 Tablestore 时序表。从工程结构看该模块规模精简、职责单一全部实现集中在以下文件中schema 与示例定义定义 Connector 与 Action 的 HOCON 配置结构、字段校验规则及 API 示例连接器资源实现实现emqx_resource行为负责客户端生命周期、健康检查、单条与批量写入动作注册 与 连接器注册向桥接框架声明动作类型与连接器类型单元测试覆盖启动、健康检查、单条/批量查询的数据构造逻辑。模块依赖ots_erlTablestore 的 Erlang 客户端以及伞形仓库内的emqx_connector、emqx_resource、emqx_bridge等基础应用其应用版本在 mix.exs 中声明。二、架构理解Connector Action 双层模型EMQX 5.x 的数据集成采用连接器 动作的两层设计Tablestore 桥接完整遵循这一模型Connector连接器描述如何连接 Tablestore。它持有实例地址、鉴权凭据、连接池等底层资源负责客户端启停与健康检查不关心具体业务数据Action动作描述把什么样的数据写到哪里。它引用某个已创建的连接器定义目标表名、measurement、tags、fields 等写入参数并支持占位符模板。在代码层面这两个角色的注册方式如下见 mix.exsenv: [ emqx_action_info_modules: [:emqx_bridge_tablestore_action_info], emqx_connector_info_modules: [:emqx_bridge_tablestore_connector_info] ]其中 emqx_bridge_tablestore_connector_info.erl 声明了bridge_types() - [tablestore, tablestore_timeseries]即该连接器同时注册了通用类型tablestore与面向时序模型的专用类型tablestore_timeseries而 emqx_bridge_tablestore_action_info.erl 则将动作类型tablestore与连接器类型tablestore绑定两者共用同一 schema 模块emqx_bridge_tablestore。连接器本身是一个标准的emqx_resource资源见 emqx_bridge_tablestore_connector.erl实现了on_start/2、on_stop/2、on_query/3、on_batch_query/3、on_get_status/2等回调从而无缝接入 EMQX 的资源管理与指标监控体系。三、创建 Tablestore Connector连接配置与健康检查3.1 连接器参数说明Connector 的 HOCON 字段定义在 emqx_bridge_tablestore.erl 的 fields(connector) 中结合 rel/i18n/emqx_bridge_tablestore.hocon 中的文案各参数含义如下参数类型必填默认值说明namestring是—连接器名称全局唯一enableboolean否true是否启用该连接器endpointstring是—Tablestore 服务地址例如https://myinstance.cn-hangzhou.ots.aliyuncs.cominstance_namestring是—Tablestore 实例名称access_key_idstring是—阿里云 AccessKey ID敏感字段schema 中标记为sensitive示例形如NTS**********************access_key_secretstring是—阿里云 AccessKey Secret敏感字段示例形如7NR2****************************************pool_sizeinteger否8与 Tablestore 交互的底层连接池大小probe_table_namestring否—用于健康检查探测的时序表名未设置时健康检查退化为列出全部时序表storage_model_typeenum否timeseries存储模型类型见下文说明SSL 相关字段—否—由emqx_connector_schema_lib:ssl_fields()附加的 TLS 配置项其中access_key_id与access_key_secret在 schema 中通过emqx_schema_secret:mk/1声明见 emqx_bridge_tablestore.erl属于敏感配置EMQX 会对其进行脱敏处理。storage_model_type在 storage_model_type_field/0 中定义为enum([timeseries])即当前仓库版本面向 Tablestore时序模型timeseries场景这也是本桥接的主要使用方式。3.2 连接器配置示例仓库在 bridge_v2_examples/1 与 connector_examples/1 中提供了可直接参考的官方示例Connector 部分如下name tablestore_connector enable true endpoint https://myinstance.cn-hangzhou.ots.aliyuncs.com storage_model_type timeseries instance_name myinstance access_key_id ****** access_key_secret ******3.3 启动与健康检查原理从 on_start/2 的实现可以看到连接器的启动流程从配置中取出instance_name、endpoint、pool_size组装基础参数access_key_id、access_key_secret作为凭据参数一起传给ots_ts_client:start/1启动 Tablestore 时序客户端启动成功后立即执行一次连通性探测probe_ots/2若配置了probe_table_name则调用ots_ts_client:describe_table/2描述该表若未配置则退化为调用ots_ts_client:list_tables/1列出全部表探测成功则返回{ok, State}并保存client_ref、probe_table_name、ots_opts及空通道表channels #{}探测失败则停止客户端并返回{error, Reason}。健康检查on_get_status/2同样依赖探测结果探测成功返回connected失败返回connecting。这意味着即使网络抖动EMQX 也能通过周期性的健康检查及时感知连接状态并将其展示在 Dashboard 与 API 中。相关的启动失败、探测表缺失、退化探测等分支均有对应的单元测试覆盖见 emqx_bridge_tablestore_connector_tests.erl 中的start_connector_failure_test_、start_connector_with_missing_probe_table_test_等用例。四、创建 Tablestore Action时序写入参数详解4.1 动作参数说明Action 的字段定义在 fields(tablestore_action) 与 fields(action_parameters) 中参数类型必填默认值说明storage_model_typeenum否timeseries存储模型类型当前支持timeseriestable_namestring模板是—目标时序表名可为静态值或占位符如${payload.table_name}measurementstring模板是—时序数据的 measurement 名称可为静态值或占位符如${payload.measurement}data_sourcestring模板否空数据源标识可为静态值或占位符如${payload.data_source}timestampinteger 或 string模板否写入时系统时间字段的微秒时间戳可为静态值或占位符如${payload.microsecond_timestamp}详见下文时间戳语义tagsmap模板否{}时序标签key 与 value 均支持静态值或占位符fields数组是非空—时序字段列表每一项为{column, value, isint?, isbinary?}schema 通过non_empty_list/1校验器拒绝空数组meta_update_modelenum否MUM_NORMAL时序元数据更新模式取值MUM_NORMAL/MUM_IGNORE其中fields数组内每个元素的子字段fields(tablestore_fields)为子字段类型必填说明columnstring模板是字段列名可为静态值或占位符${payload.column}valueboolean / number / string模板是字段值可为静态值或占位符${payload.value}isintboolean / string模板否是否将数值按整数写入。默认false即整数按浮点数写入isbinaryboolean / string模板否是否将二进制值按 binary 类型写入。默认false即二进制值按字符串写入isint与isbinary的语义在 rel/i18n/emqx_bridge_tablestore.hocon 中有明确说明它们同样支持占位符例如${payload.is_int}、${payload.is_binary}因此可以在消息层面动态决定字段的类型映射。4.2 动作配置示例仓库提供的官方示例action_parameters_example/1如下parameters { storage_model_type timeseries table_name ${data_source} measurement ${measurement} tags { tag1 ${tag1} tag2 ${tag2} } fields [ { column ${column} value ${value} isint true } ] meta_update_model MUM_IGNORE }动作的批量写入参数fields(action_resource_opts)默认值为batch_size 100、batch_time 100ms即每条消息到达动作时并不立即落库而是进入缓冲攒满 100 条或等待 100ms 后批量写入这与 EMQX 桥接框架中action_resource_opts的通用语义一致见 emqx_bridge_v2_schema.erl。五、占位符与模板渲染机制Tablestore 桥接的动作参数广泛支持占位符模板这是其与规则引擎联动、实现消息字段动态映射到时序数据的关键机制。在连接器端on_add_channel/4 会为每个动作通道预编译参数凡是包含${前缀的二进制串都会被maybe_preproc/1转换成{tmpl_tokens, Tokens}模板令牌通过emqx_placeholder:preproc_tmpl/1预编译静态值则原样保留tags与fields中的每个 key/value 也会被逐项预编译。当消息到达时渲染过程如下table_name、measurement、data_source、timestamp直接通过render_tmpl/2渲染tags通过render_tags/2遍历渲染只有 key 与 value 都渲染成功的项才会进入最终标签集合渲染失败返回undefined的项被静默丢弃fields通过render_fields/2遍历渲染column 与 value 均渲染成功的项保留同时根据isint、isbinary是否渲染出值来决定是否携带对应类型选项若table_name或measurement渲染结果为undefined则抛出{bad_ots_data, no_table_name}/{bad_ots_data, no_measurement}最终被捕获并转换为{error, {unrecoverable_error, _}}返回——这保证了下游永远不会收到残缺的写入请求。占位符的取值来自消息本身典型形如${payload.field_name}、${tag_name}等具体可用的占位符集合由上游规则引擎 SQL 的 SELECT 字段决定。单元测试 on_query_test_ 验证了完整渲染链路例如tags中的${tag1} ${tag1_value}最终渲染为#{tag1 : tag1_value}fields中带isint true的整数项渲染为{int_field, 123, #{isint : true}}。六、数据写入语义时间戳、元数据更新与批量写入6.1 时间戳语义timestamp参数的解析逻辑在 do_mk_tablestore_data_row/3 中实现规则如下未配置、渲染结果为undefined或渲染结果为字符串now/NOW时使用写入时刻的系统微秒时间戳os:system_time(microsecond)渲染结果为整数直接作为时间戳使用其他情况抛出{bad_ots_data, {bad_timestamp, Ts}}拒绝写入。因此在动作中不配置timestamp即可获得写入时自动打点的行为而若消息携带业务时间则可通过${payload.microsecond_timestamp}之类的占位符传入微秒级时间戳。6.2 元数据更新模型meta_update_model控制 Tablestore 时序元数据即 measurement、tags 等描述信息的更新行为见 rel/i18n/emqx_bridge_tablestore.hoconMUM_NORMAL默认正常模式写入时若时序元数据不存在则自动创建MUM_IGNORE忽略模式写入时不尝试创建时序元数据。该值会原样透传给底层ots_ts_client:put/2在 mk_tablestore_data/3 中以meta_update_mode字段进入请求体适合对元数据创建行为有严格管控的场景可以避免高频写入带来的元数据频繁创建开销。6.3 单条与批量写入单条写入on_query/3针对一条消息渲染出一行数据后调用ots_ts_client:put/2写入批量写入on_batch_query/3接收一批{ChannelId, Message}先按渲染出的table_name对消息分组再为每个表名构造一批rows_data逐表调用put/2见 mk_tablestore_batch_data/2。批量中的任何一次put失败都会汇总为{error, {unrecoverable_error, Errors}}返回。写入的数据行结构为#{measurement, data_source, tags, fields, time}其中data_source未配置时为空字符串trans_data_source/1将undefined转为。批量写入的测试用例 on_batch_query_test_ 验证了 3 条消息在同一表名下被合并为 3 行数据一次提交。七、通过规则引擎将消息路由到 Tablestore桥接在 EMQX 中属于数据集成能力其典型用法是在规则引擎中创建一条规则用 SQL 筛选出需要入库的消息然后将匹配结果作为动作Action发送到 Tablestore 桥接。大致的集成流程为创建 Tablestore Connector上文的连接器配置创建 Tablestore Action引用该连接器并配置table_name、measurement、tags、fields等写入参数占位符与规则 SQL 的字段一一对应在规则引擎中编写 SQL例如SELECT payload.*, clientid FROM t/#将需要的字段投射到消息中将规则的动作指向上面创建的 Tablestore Action。当消息命中规则时EMQX 会将其转换为动作所需的输入携带 SQL 中 SELECT 出的字段连接器中的模板渲染逻辑随即把字段填充进时序数据行并写入 Tablestore。规则引擎的详细介绍与 SQL 语法可参考仓库内的 emqx_rule_engine 应用及其 README桥接作为规则动作的使用方式与 EMQX 数据集成文档保持一致配置入口同时支持 Dashboard 图形化操作与 HTTP API 两种方式。八、桥接管理 HTTP APIemqx_bridge_tablestore模块为桥接管理提供了完整的 HTTP API 支持涵盖创建、更新、查询、停止/重启与列表等操作。这些 API 的请求/响应 schema 定义在 emqx_bridge_tablestore.erl 的fields/1中连接器管理get_connector、put_connector、post_connector见fields(Field) when Field get_connector; ...动作管理get_bridge_v2、post_bridge_v2、put_bridge_v2动作以bridge_v2命名空间暴露。调用这些 API 时请求体可直接复用仓库提供的官方示例GET/POST/PUT三种方法对应的配置结构一致例如创建连接器时使用 connector_examples/1 中的tablestore_timeseries示例创建动作时使用 bridge_v2_examples/1 中的tablestore_timeseries示例。API 的具体路径与鉴权方式遵循 EMQX Dashboard/管理 API 的通用规范可在 Dashboard 的数据集成页面或 API 文档中找到对应端点。九、测试与验证仓库为该桥接提供了较为完整的单元测试见 emqx_bridge_tablestore_connector_tests.erl其覆盖场景可直接作为接入前的自查清单测试用例验证内容start_connector_test_连接器正常启动客户端引用返回、endpoint 与 pool_size 正确透传start_connector_failure_test_Tablestore 返回OTSParameterInvalid时启动失败start_connector_with_missing_probe_table_test_探测表不存在OTSObjectNotExist时启动失败start_connector_list_tables_fallback_test_未配置探测表时退化调用list_tables成功启动on_get_status_describe_probe_test_配置探测表时describe_table成功 →connected失败/超时 →connectingon_get_status_list_tables_fallback_test_未配置探测表时list_tables成功 →connected失败 →connectingon_query_test_单条消息的完整渲染表名、measurement、时间戳、tags、fields含 int/float/bool/binary 类型选项on_batch_query_test_批量消息按表名分组合并写入这些测试通过 meck 模拟ots_ts_client行为无需真实 Tablestore 实例即可验证桥接的数据构造逻辑对二次开发或排查问题很有参考价值。十、小结emqx_bridge_tablestore以连接器 动作的标准数据集成模型将 EMQX 的规则引擎能力与阿里云表格存储的时序模型无缝衔接连接器层负责凭据、连接池与基于探测表的健康检查动作层通过丰富的占位符模板把消息字段动态映射为时序表的 measurement、tags、fields 与时间戳并以MUM_NORMAL/MUM_IGNORE控制元数据更新行为配合默认batch_size 100、batch_time 100ms的批量写入获得较高的吞吐表现。无论你是在搭建车联网、工业 IoT 还是其他时序数据采集场景都可以参考本文的配置示例与源码导读快速完成从 MQTT 消息到 Tablestore 时序数据的落库链路。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐SeaTunnel Tablestore Sink 连接器实战指南将海量数据写入阿里云表格存储SeaTunnel Tablestore Sink 连接器实战指南将海量数据写入阿里云表格存储 SeaTunnel 的 Tablestore Sink 连接器数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Tablestore Sink 连接器实战指南将数据高效写入阿里云 TablestoreSeaTunnel Tablestore Sink 连接器实战指南将数据高效写入阿里云 Tablestore 本指南围绕 SeaTunnel 的 Tables数据集成ETL大数据批处理流处理变更数据捕获LlamaIndex 阿里云表格存储TablestoreDocumentStore 集成实战基于 TablestoreKVStore 的文档存储方案LlamaIndex 阿里云表格存储TablestoreDocumentStore 集成实战基于 TablestoreKVStore 的文档存储方案 导读人工智能RAG大模型上一篇AMD GLM-4.7-MXFP4性能评测99.68%精度恢复的量化奇迹下一篇Eclipse Theia AI Registry 扩展一站式浏览与安装经批准的 MCP 服务器创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考