- 后端
- 物联网
- 消息队列
- 通信
【免费下载链接】emqx
The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles
本文基于 EMQX 开源仓库
changes/ee/fix-17216.en.md变更说明,结合 emqx_bridge_pgsql 与 emqx_postgresql 的源码与测试用例,深入解析该修复的背景、触发场景、修复机制与验证方式,帮助读者理解 EMQX 数据集成中 SQL 参数类型映射的错误处理链路,并在实际使用中规避此类数据问题。
一、修复内容速览
该变更(对应内部工单fix-17216)针对 EMQX 的TimescaleDB 与 PostgreSQL 桥接动作(bridge action)修复了一个故障模式:
当一条带引号的 JSON 数字字符串(例如
"27.1")被映射写入FLOAT类型列时,桥接动作不再导致数据库连接进程崩溃,而是返回一个结构化的参数错误(structured bad parameter error)。
修复的核心价值在于:
- 错误可观测:错误以结构化字段(
reason、type、index、value)呈现,便于日志采集、告警与排障; - 连接保持可用:数据库连接进程不再因类型编码失败而崩溃,避免了连接池内连接被逐个击穿、动作持续不可用的连锁故障;
- 行为可预期:该错误被归类为不可恢复错误(unrecoverable error),规则引擎可以据此明确判定该次写入失败,而不会产生含糊的崩溃日志。
这一修复同时作用于PostgreSQL、TimescaleDB 以及复用同一套实现的 Matrix 桥接(详见后文源码分析)。
二、触发场景:JSON 数字字符串与 FLOAT 列的映射冲突
2.1 场景描述
在 EMQX 数据集成(Data Integration)中,用户通常通过规则引擎(Rules)订阅 MQTT 主题,将消息负载(payload)中的字段通过 SQL 模板映射到目标数据库表列。
当出现以下条件组合时,便会触发该修复对应的故障:
- 消息 payload 中某字段的值是带引号的数字字符串,例如 JSON 文本
{"temp":"27.1"}中的temp值为字符串"27.1"(而非数字27.1); - 桥接动作的 SQL 模板将该字段直接映射到目标表中类型为
FLOAT(底层为 PostgreSQLfloat8)的列; - 桥接使用epgsql 驱动(EMQX PostgreSQL/TimescaleDB 桥接所依赖的 Erlang PostgreSQL 客户端)执行带参数绑定的写入。
2.2 为什么这是一个高频坑
IoT 场景下,传感器数据经常以 JSON 文本形式发布,而 JSON 数字在消息序列化过程中可能被保留为字符串形式(例如网关固件、历史数据回放、不同厂商协议转换后的数据),因此"27.1"这类“看起来是数字、实际上是字符串”的值非常常见。与此同时,TimescaleDB 是构建在 PostgreSQL 之上的时序数据库扩展,温度、湿度、电压等时序指标列普遍定义为FLOAT/DOUBLE PRECISION。二者相遇时,类型不匹配几乎是必然的。
2.3 修复前的故障表现
在修复之前,当 epgsql 驱动尝试将一个二进制字符串(Erlang binary,如<<"27.1">>)编码为float8参数时,编码失败会向上抛出异常,最终导致承载数据库连接的 worker 进程崩溃。对于使用连接池(pool)的桥接而言,每次写入失败都可能击穿池内一个连接,造成:
- 连接进程反复崩溃、重建,产生大量噪声日志;
- 该桥接动作的健康状态恶化,甚至触发连接池整体不可用;
- 排障困难:崩溃日志往往只暴露驱动层面的异常,难以定位到具体是哪条消息、哪个字段、哪个值导致的失败。
三、修复机制:从崩溃到结构化错误的结构性变化
3.1 错误处理链路源码佐证
修复的核心逻辑位于 apps/emqx_postgresql/src/emqx_postgresql.erl 的handle_result/1函数(约 L878-L880):
handle_result({error, #{reason := bad_param} = Context}) -> ?tp("postgres_bad_param_error", #{context => Context}), {error, {unrecoverable_error, Context}};这里的关键在于:
bad_param结构化错误:当查询结果中携带reason := bad_param上下文时,会发布postgres_bad_param_error追踪事件(snabbkaffe trace point),并返回{error, {unrecoverable_error, Context}};Context完整透传:上下文中的type(参数类型,如float8)、index(参数在 SQL 中的位置)、value(触发错误的原始值)等信息全部保留,供上层消费;- 不可恢复错误语义:
unrecoverable_error明确告知规则引擎:该次写入失败是确定性的、由数据本身导致,重试也不会成功,从而避免无效重试放大故障。
3.2 错误翻译与导出
在错误最终返回给上层之前,会经过translate_to_log_context/1与export_error/1处理(emqx_postgresql.erl 约 L912-L946):
translate_to_log_context(#error{} = Reason) -> #{ driver_severity => Severity, driver_error_codename => Codename, driver_error_code => Code, driver_error_message => ..., driver_error_extra => Extra };这保证了:
- 日志侧:错误以结构化的
driver_error_*字段写入日志,可被日志系统直接检索、聚合; - API 侧:
export_error/1将severity、error_codename、error_code等关键字段导出到桥接动作的查询结果中,用户可以通过规则引擎、Dashboard 或 API 查询到明确的失败原因。
3.3 连接进程为何不再崩溃
从整体调用链看(emqx_postgresql.erl 的on_sql_query/6约 L496-L546):
- 查询执行通过
ecpool:pick_and_do/4在连接池 worker 上执行,结果以返回值形式返回,而非让异常逃逸出 worker 进程; on_sql_query/6对{error, Reason}结果统一进行错误翻译与日志记录,再经handle_result/1分类处理;- 即使是驱动抛出的异常(如
error:function_clause),也会被try...catch捕获(约 L535-L546),并转换为{error, {unrecoverable_error, invalid_request}},而不会直接击穿 worker 进程。
因此,bad_param类错误在修复后走的是完整的结构化错误返回路径,worker 进程存活,连接池状态稳定。
四、修复验证:测试用例逐条解读
该修复有对应的回归测试覆盖,位于 apps/emqx_bridge_pgsql/test/emqx_bridge_pgsql_SUITE.erl。
4.1 核心回归用例t_bad_float_param(约 L910-L960)
测试流程还原了真实故障链路:
t_bad_float_param(TCConfig) -> Conn = connect_direct_pgsql(TCConfig), {ok, _, _} = epgsql:squery(Conn, <<"ALTER TABLE mqtt_test ADD COLUMN temp FLOAT">>), ok = epgsql:close(Conn), {201, _} = create_connector_api(TCConfig, #{}), {201, _} = create_action_api(TCConfig, #{ <<"parameters">> => #{ <<"sql">> => << "INSERT INTO mqtt_test(payload, temp) " "VALUES (${payload}, ${temp})" >> } }), ... Payload = <<"{\"temp\":\"27.1\"}">>, ... ?assertMatch( {_, {ok, #{ context := #{ reason := bad_param, type := float8, index := 1, value := <<"27.1">> } }}}, ?wait_async_action( emqtt:publish(C, RuleTopic, Payload), #{?snk_kind := "postgres_bad_param_error"} ) ),要点解析:
- 构造场景:给测试表
mqtt_test增加temp FLOAT列,并创建将 payload 的temp字段直接写入该列的桥接动作; - 发布带引号数字:向规则主题发布
{"temp":"27.1"},其中temp是字符串"27.1"; - 断言结构化错误:期望返回的上下文中
reason := bad_param、type := float8、index := 1、value := <<"27.1">>,且伴随postgres_bad_param_error追踪事件; - 指标断言:随后验证桥接动作指标——
matched_get为 1(消息被匹配)、failed_get为 1(写入失败)、success_get为 0,并且表内行数仍为 0,证明失败被精确记录、没有产生脏数据。
4.2 配套测试:同类参数错误的统一处理
同一测试套件还包含对时间戳类型的同类回归测试t_bad_datetime_param(约 L880-L908):将非法时间戳值写入timestamp列时,同样返回reason := bad_param, type := timestamp, index := 1, value := ...的结构化错误。这证明bad_param结构化错误是一个通用机制,覆盖 float8、timestamp 等多种参数类型,而非针对单一类型的特判。
此外,t_bad_sql_parameter(约 L607-L638)验证了驱动层参数编码失败的兜底路径:当参数本身无法编码时,批量(batch)模式返回{error, {unrecoverable_error, invalid_request}},同步模式返回{error, {unrecoverable_error, _}},同样不会崩溃连接进程。
4.3 测试参数化与多后端覆盖
t_bad_float_param的矩阵定义:
t_bad_float_param() -> [{matrix, true}]. t_bad_float_param(matrix) -> [[?timescale, ?sync, ?without_batch]];即该用例通过matrix参数化框架,在TimescaleDB、同步模式、非批量组合下运行,直接对应修复工单标题中的 TimescaleDB 场景。
五、为什么 Timescale、PostgreSQL、Matrix 一起修复
5.1 共享的 schema 定义
TimescaleDB 桥接的 HOCON schema 直接复用了 PostgreSQL 桥接的 action schema。在 apps/emqx_bridge_timescale/src/emqx_bridge_timescale.erl 中:
fields("post") -> emqx_bridge_pgsql:fields("post", ?ACTION_TYPE, "config"); fields("put_bridge_v2") -> emqx_bridge_pgsql:fields(pgsql_action);Timescale 桥接的 connector schema 则复用 apps/emqx_postgresql/src/schema/emqx_postgresql_connector_schema.erl 的定义。也就是说,TimescaleDB 桥接在配置结构上就是“PostgreSQL 桥接 + Timescale 类型名”。
5.2 共享的底层驱动实现
所有基于 PostgreSQL 的桥接(pgsql、timescale、matrix、以及复用同一 schema 的 redshift、cockroachdb、alloydb 等)最终都通过 apps/emqx_postgresql/src/emqx_postgresql.erl 中的资源回调(on_query、on_batch_query、on_start、on_get_status等)与 epgsql 驱动交互。因此,handle_result/1中bad_param的处理逻辑被所有相关桥接共享——一处修复,全线生效。
5.3 桥接动作的参数解析与模板渲染
桥接动作执行时,SQL 模板经 emqx_template_sql 解析为带参数的预编译语句模板(见 emqx_postgresql.erl 的parse_sql_template/2与render_prepare_sql_row/2),payload 字段经 JSON 语义(emqx_jsonish)渲染为参数行。参数最终由 epgsql 驱动按目标列类型编码。当编码器收到无法编码为float8的二进制字符串时,即触发本次修复所处理的bad_param错误路径。
六、对使用者的实战建议
6.1 如何避免触发该错误
- 在规则中做类型转换:在 SQL 模板或规则 SQL 中使用类型转换表达式,将字符串数字显式转为数值后再写入,例如在 SQL 模板中写
(${temp} :: float),或使用 EMQX 规则引擎的 SQL 函数(如float()之类)完成转换; - 数据侧治理:在网关或边缘侧规范 payload 中数字字段的 JSON 类型,避免数字被序列化为字符串;
- 合理设计表结构:若数据源数字字段类型不稳定,可将目标列定义为
TEXT或JSONB,在查询时再做转换,或使用 TimescaleDB 的连续聚合/物化视图承接后续计算。
6.2 错误发生后的排查路径
- 在桥接动作指标中查看
failed计数(对应测试中的failed_get),确认失败是否由数据写入引起; - 在日志中检索
postgres_bad_param_error或postgresql_connector_do_sql_query_failed相关条目,其中driver_error_codename、driver_error_message字段会给出驱动层面的具体原因; - 通过 Dashboard 或 API 查询该动作的
status与status_reason,确认连接健康状态未受影响。
6.3 版本适用范围说明
- 本文所述行为以当前开源仓库(apps/emqx_bridge_pgsql、apps/emqx_bridge_timescale、apps/emqx_postgresql)的实现为准;
- 涉及的具体配置项包括桥接动作的
parameters.sql模板、连接的disable_prepared_statements(默认false,开启预编译语句)、以及动作的resource_opts(如batch_size默认 100、batch_time默认100ms,见 emqx_bridge_pgsql.erl L77-L81); - 默认 SQL 模板为
INSERT INTO t_mqtt_msg(msgid, topic, qos, payload, arrived) VALUES (${id}, ${topic}, ${qos}, ${payload}, TO_TIMESTAMP((${timestamp} :: bigint)/1000))(见 emqx_bridge_pgsql.erl L134-L138)。
七、小结
fix-17216的核心成就在于把“带引号 JSON 数字字符串映射 FLOAT 列导致数据库连接进程崩溃”这一隐蔽故障,收敛为“结构化、可观测、不可恢复的参数错误”。修复后的行为让数据集成链路对坏数据具备了更强的韧性:单条坏数据只会导致该条写入失败并被明确标记,而不会拖垮连接池、污染健康状态或产生无法定位的崩溃日志。对于以 TimescaleDB/PostgreSQL 作为时序数据落库后端的 EMQX 用户而言,理解这一修复机制有助于更好地设计规则 SQL、规避类型陷阱,并在故障发生时快速定位根因。
- 后端
- 物联网
- 消息队列
- 通信
【免费下载链接】emqx
The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles
相关推荐
EMQX 插件配置 API 修复:根 JSON 值类型错误时返回可读校验信息而非 500
EMQX 插件配置 API 修复:根 JSON 值类型错误时返回可读校验信息而非 500 导读 本篇文章基于 EMQX 开源仓库中的变更记录 fix 18153
后端物联网消息队列通信EMQX 集成 TimescaleDB 数据桥接实战:基于 PostgreSQL 扩展的时序数据写入方案
EMQX 集成 TimescaleDB 数据桥接实战:基于 PostgreSQL 扩展的时序数据写入方案 导读 本文围绕 EMQX 开源仓库中的 emqx_br
后端物联网消息队列通信PHPStan `new.dateInterval` 错误详解:DateInterval 构造函数的非法时长字符串检测与修复
PHPStan new.dateInterval 错误详解:DateInterval 构造函数的非法时长字符串检测与修复 PHPStan(PHP Static
开发工具代码质量静态分析
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考