- 后端
- 物联网
- 消息队列
- 通信
【免费下载链接】emqx
The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles
导读
本文围绕 EMQX 仓库中 Oracle 数据库桥接(emqx_bridge_oracle/emqx_oracle)的一项关键修复展开:prepare/status 检查改为"只解析、不执行"用户 Action SQL,并拒绝顶层 DDL/DCL/TCL 语句,同时改进超过 4000 字节文本 payload 的绑定支持。读完本文,你将理解 EMQX 在创建 Oracle Action 时如何安全地校验 SQL 模板(既不触发数据变更,又能提前发现表不存在等错误),以及如何让大体积 MQTT 消息文本顺利写入 Oracle NCLOB 列。全文以 changes/ee/fix-17605.en.md 为骨架,结合仓库源码与测试用例给出可验证的实现细节。
一、问题背景:Oracle Action 的 SQL 模板与 Prepare 阶段
EMQX 规则引擎通过数据桥接将 MQTT 消息写入 Oracle 时,Action 的核心配置是 SQL 模板(sql字段)。该模板使用${topic}、${payload}等占位符,运行时由emqx_placeholder:preproc_sql/2将模板预处理为带命名绑定的预编译语句,再通过 jamdb_oracle 驱动以绑定参数方式执行(见 emqx_oracle.erl 的parse_prepare_sql与init_prepare)。
修复前存在两个痛点:
- prepare/status 检查会真正执行用户 SQL:为验证模板可解析、目标表存在,旧实现可能直接执行语句,带来副作用风险;如果模板被误写成
CREATE TABLE、DROP TABLE等语句,一旦"准备"阶段执行,将直接修改数据库结构。 - Oracle 文本绑定存在 4000 字节上限问题:Oracle 的
VARCHAR2绑定默认上限为 4000 字节,超过该长度的文本 payload(例如较大的 JSON 消息体)在绑定阶段会失败,导致大消息无法入库。
本次修复(issue fix-17605)同时解决了这两类问题,其成果集中在两个层面:
emqx_oracle连接器驱动层:新增"解析式预检 + 语句类型白名单"机制;emqx_bridge_oracle桥接测试层:新增针对"不执行用户 SQL""拒绝 DDL""大 payload"的 CT 用例。
二、核心修复一:prepare/status 检查改为"只解析、不执行"
2.1 用dbms_sql.parse匿名块替代直接执行
在 emqx_oracle.erl 中,所有 prepare 路径(连接池建立、channel 添加、健康检查)都收敛到同一个入口:
check_if_table_exists(Conn, SQL, _Tokens0) -> case sql_has_parse_side_effect(SQL) of true -> {error, unsupported_sql_statement}; false -> check_if_sql_parseable(Conn, SQL, 1) end.其中check_if_sql_parseable/3(源码位置)构造了一段 PL/SQL 匿名块,借助 Oracle 内置包dbms_sql的parse过程只做语法解析而不执行语句:
declare c integer; begin c := dbms_sql.open_cursor; dbms_sql.parse(c, :1, dbms_sql.native); dbms_sql.close_cursor(c); exception when others then if dbms_sql.is_open(c) then dbms_sql.close_cursor(c); end if; raise; end;调用时把用户 SQL 原文作为唯一绑定参数:1传入({ParseSQL, [binary_to_list(SQL)]})。这样:
- 解析成功:返回
{ok, [{proc_result, 0, _}]},即ok,模板可安全使用; - 解析失败:
proc_result返回 Oracle 错误码(如 904/942/1013 等),进入错误归类逻辑(详见第四节)。
从实现看,这套预检在
prepare_sql_to_conn/4(源码)中被每个连接调用,并在连接池注册为 reconnect callback,因此连接重建后依然会重新校验。
2.2 测试如何证明"没有执行"
仓库用两条测试用例从正反两面验证"只解析、不执行":
- 正面验证(执行了就会留下痕迹):t_prepare_does_not_execute_user_sql 先在 Oracle 上创建一个带
BEFORE INSERT触发器(sql_create_probe_trigger,见 SUITE 第 310-318 行)和一张审计表mqtt_probe_audit:触发器在向mqtt_test表 INSERT 时,以自治事务向审计表写入一条'prepare'标记。随后正常创建 connector 与 action 并触发 status 检查,最后断言mqtt_probe_audit行数为0—— 证明 prepare 阶段从未真正执行 INSERT。 - 反向验证(拒绝 DDL 且不建表):t_prepare_rejects_ddl_without_executing 将
CREATE TABLE mqtt_prepare_probe (id NUMBER)作为 Action SQL 提交,断言该表并未被创建(count_user_table(..., "MQTT_PREPARE_PROBE")为 0),同时轮询 Action 状态接口,最终status_reason中应包含unsupported_sql_statement。
三、核心修复二:拒绝顶层 DDL/DCL/TCL 语句
3.1 词法级语句类型检查
check_if_table_exists先调用sql_has_parse_side_effect/1(源码)做一次纯本地词法检查:提取 SQL 的首个有效 token,与一组"有解析副作用"的语句关键字比对,命中即直接返回{error, unsupported_sql_statement},根本不会发起网络请求:
| 分类 | 被拒绝的关键字 |
|---|---|
| DDL(数据定义语言) | alter、create、drop、truncate、rename、comment、flashback、purge、analyze、associate、disassociate、audit、noaudit |
| DCL(数据控制语言) | grant、revoke、lock |
| TCL(事务控制语言) | commit、rollback、savepoint、set |
| 其他 | explain(EXPLAIN PLAN) |
对应的错误消息宏定义在 emqx_oracle.erl 第 14-20 行:
-define(UNSUPPORTED_SQL_STATEMENT_MSG, "unsupported_sql_statement: DDL, DCL, and transaction control statements are not supported " "in Oracle Action SQL templates." ).3.2 token 提取对注释与空白是健壮的
first_sql_token/1及其辅助函数(源码)在取首个 token 前会跳过:
- 空白字符(空格、
\t、\n、\r、\f、\v); - 行注释
-- ...; - 块注释
/* ... */。
token 统一转小写比对,因此CREATE、Create、create都会被识别。仓库内置的 eunit 测试first_sql_token_test_(第 788-798 行)覆盖了这些场景,例如"-- comment\nSELECT 1 FROM dual"会正确解析出select,"/* unterminated comment"解析为空 token。
3.3 白名单之外:合法 Action SQL 仍然放行
sql_has_parse_side_effect_test_(第 742-786 行)明确列出被放行的常规语句类型,即 Action 模板的合法范围:
INSERT INTO ... VALUES(:1, :2)(EMQX 默认模板即 INSERT);UPDATE ... SET ... WHERE ...;DELETE FROM ... WHERE ...;MERGE INTO ... USING dual ON ...;SELECT ... WHERE ...与WITH ... SELECT ...;- 匿名 PL/SQL 块
BEGIN NULL; END;(用于存储过程调用场景,对应 SUITE 中的sql_stored_procedure_template)。
四、错误归类与健康状态机:unsupported_sql_statement / undefined_table
预检的错误需要映射为桥接的健康状态,相关逻辑集中在 on_get_channel_status 与do_check_prepares/3(第 373-414 行):
| 预检结果 | Action 健康状态 | 状态说明 |
|---|---|---|
ok | connected | 模板可解析、目标表存在 |
{error, undefined_table} | disconnected+{unhealthy_target, "Oracle table is invalid. Please check if the table exists..."} | 表/视图不存在或标识符无效 |
{error, unsupported_sql_statement} | disconnected+{unhealthy_target, "unsupported_sql_statement: DDL, DCL, and transaction control statements are not supported..."} | 顶层语句类型不被允许 |
| 其他错误 | connecting | 暂时性失败,等待下一次健康检查 |
do_check_prepares会遍历连接池中的每个 worker 连接,逐一用check_if_table_exists验证,因此即使部分连接异常也能被捕获。
check_if_sql_parseable/3对 Oracle 返回码的归类规则(第 578-599 行)与handle_parse_error_description/1(第 601-611 行):
ORA-00904: invalid identifier→{error, undefined_table}(列名写错,如测试中的retainx);ORA-00942: table or view does not exist→{error, undefined_table};ORA-01013: user requested cancel of current operation→ 视为可重试,重试一次;- 其他 ORA 错误 → 原样返回描述,作为其他错误处理。
错误码解析使用正则ORA-([0-9]+)从描述文本中提取(oracle_error_codes/1,第 613-624 行)。此外,针对noproc(连接进程重启竞态,见源码中的Note [jamdb oracle race condition])与socket closed两类瞬时故障,check_if_sql_parseable会带重试计数自动重试一次。
五、大文本 payload(超过 4000 字节)的绑定支持
Oracle 传统VARCHAR2绑定参数存在 4000 字节上限,而 MQTT 消息 payload 完全可能超过该长度。本次修复同时改善了文本型 payload 的写入支持,条件与效果如下:
- 条件:payload 占位符(
${payload})位于 SQL 模板的最后一个绑定参数位置; - 效果:超过 4000 字节的文本 payload 可以正常绑定并写入。
仓库在桥接层将 payload 列设计为NCLOB(见测试建表语句 第 194-195 行payload NCLOB,以及 第 242-251 行 的t_mqtt_msgs表),NCLOB 用于承载大文本。
测试用例给出可直接验证的证据链:
- t_probe_with_large_value 使用模板
"INSERT INTO mqtt_test(topic, msgid, payload, retain) VALUES (${topic}, ${id}, ${payload}, ${id})"(${payload}为最后一个绑定参数)执行 probe,预期返回 204(成功),证明大值在预检阶段即可通过; - t_eec_1322_large_payload_after_small_payload 先写入小 payload(
{"msg":"heelo"}),再写入large_json_payload()—— 即binary:copy(<<"a">>, 5000)生成的 5000 字节文本经 JSON 编码后的 payload(第 513-514 行),断言两行均成功落库、Action 状态保持connected,且 trace 中无oracle_connector_query_return错误事件。
六、配置与使用:如何应用这些修复
6.1 Action 的 SQL 模板配置
Oracle Action 的sql参数在 emqx_bridge_oracle.erl 中定义,类型为emqx_schema:template(),默认值:
insert into t_mqtt_msgs(msgid, topic, qos, payload) values (${id}, ${topic}, ${qos}, ${payload})schema 同时给出 Action 级批量参数默认值(第 131-135 行):batch_size默认 100、batch_time默认100ms,与 CT 用例中的?with_batch矩阵(batch_size=100, batch_time=200ms)对应。
6.2 Connector 连接配置要点
连接器配置项定义于 emqx_oracle_schema.erl:
| 配置项 | 类型 | 说明 |
|---|---|---|
server | host:port | 默认端口 1521(emqx_oracle:oracle_host_options/0) |
sid | binary,可选 | Oracle SID,与service_name至少填一个 |
service_name | binary,可选 | Oracle 服务名,与sid至少填一个 |
role | normal/sysdba | 连接角色,默认normal |
username/password | string | 必填,测试默认system/oracle |
pool_size | integer | 连接池大小,默认 8(见?DEFAULT_POOL_SIZE) |
emqx_bridge_oracle的config_validator/1(第 194-203 行)会在sid与service_name都缺失时返回校验错误"neither SID nor Service Name was set",对应测试 t_no_sid_nor_service_name。
6.3 编写 Action SQL 模板的实践建议
综合本次修复,编写 Oracle Action SQL 时应遵循:
- 只使用 DML(INSERT/UPDATE/DELETE/MERGE)或匿名 PL/SQL 块,禁止在模板顶层出现
CREATE、DROP、TRUNCATE、GRANT、COMMIT等 DDL/DCL/TCL 语句 —— 否则 Action 将处于disconnected状态,status_reason显示unsupported_sql_statement; - 目标表必须在创建 Action 前建好,否则预检报
ORA-00942,状态显示目标表无效(unhealthy_target); - 超过 4000 字节的文本 payload,建议将
${payload}放在 SQL 模板的最后一个绑定参数位置,并将目标列设计为NCLOB,以获得最佳兼容性; - 模板中的占位符支持嵌套取值(如
${payload.msg}),可配合emqx_placeholder预处理;字段缺失时绑定为NULL(见 t_message_with_null_value)。
七、相关源码与测试索引
如需深入阅读,可关注以下仓库文件:
- apps/emqx_oracle/src/emqx_oracle.erl:核心实现,含
check_if_table_exists、check_if_sql_parseable、sql_has_parse_side_effect、do_check_prepares、on_get_channel_status,以及-ifdef(TEST)内的全套 eunit 用例; - apps/emqx_oracle/src/emqx_oracle_schema.erl:连接器 HOCON schema(server/sid/service_name/role 等);
- apps/emqx_bridge_oracle/src/emqx_bridge_oracle.erl:Oracle Action 与 Connector 的桥接 API schema、默认 SQL 模板与参数校验;
- apps/emqx_bridge_oracle/test/emqx_bridge_oracle_SUITE.erl:CT 测试,覆盖"prepare 不执行用户 SQL""拒绝 DDL""大 payload 顺序写入""存储过程模板""表被删除/缺失""空值/嵌套 token"等场景。
结语
fix-17605 这项修复让 EMQX 的 Oracle 数据桥接在"SQL 模板安全预检"与"大文本写入"两个维度上更加稳健:预检阶段通过dbms_sql.parse匿名块做到只解析、不执行,配合首 token 词法检查从源头拒绝 DDL/DCL/TCL,再通过 ORA 错误码归类驱动健康状态机,使问题在规则真正触发前即可被发现;而对超过 4000 字节文本 payload 的绑定支持,则让大体积 MQTT 消息能够可靠地落库到 Oracle NCLOB 列。上述行为均有源码与测试用例可查证,可直接作为排查 Oracle 桥接状态异常(unsupported_sql_statement、unhealthy_target)的参考依据。
- 后端
- 物联网
- 消息队列
- 通信
【免费下载链接】emqx
The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles
相关推荐
EMQX 5.1.0 版本深度解读:连接保活、数据桥接与安全加固的全面升级
EMQX 5.1.0 版本深度解读:连接保活、数据桥接与安全加固的全面升级 导读 本文基于当前仓库中 changes/e5.1.0.en.md https://
后端物联网消息队列通信EMQX Oracle 数据库连接器(emqx_oracle)深度解析:连接管理、SQL 模板与数据桥接实战
EMQX Oracle 数据库连接器(emqx_oracle)深度解析:连接管理、SQL 模板与数据桥接实战 导读 本文围绕 EMQX 仓库中的 Oracle
后端物联网消息队列通信PGlite 事务安全加固:事务句柄在事务关闭后被拒绝执行
PGlite 事务安全加固:事务句柄在事务关闭后被拒绝执行 导读 本篇文章围绕 PGlite 一个重要的防御性变更展开: 当事务(Transaction)已经结
数据库嵌入式数据库WebAssembly
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考