为什么 MySQL 层需要单独设计
PipeMonitor 的 API 服务在 MQTT 和 Flutter App 之间承担三件事:把上行帧落库、对外提供历史查询、保证补传场景下数据不重复。前两个看起来只是 CRUD,但只要把“DR154 断网补传”和“STM32 重启 seq 归零”这两件事放进来,数据库层就会冒出一堆边界问题。
src/db.js 把这些边界集中在一个 createDbStore 工厂里,启动时建表、建索引、清理脏数据,运行时用 ON DUPLICATE KEY UPDATE 兜住补传,查询时用 received_at 而不是设备时间戳过滤。这篇把几个关键设计拆开讲。
连接池与有限次重试
容器编排里 MySQL 容器未必比 API 容器先就绪。如果 API 启动时 MySQL 还在初始化数据目录,直接 createPool 后发查询会拿到 ECONNREFUSED。initializePool 做了有限次重试,而不是无限等待:
const DB_CONNECT_ATTEMPTS = 30;
const DB_CONNECT_RETRY_MS = 2000;
async function initializePool(pool, config, status) {
// 容器编排场景里 MySQL 可能晚于 API 就绪,这里做有限次重试。
for (let attempt = 1; attempt <= DB_CONNECT_ATTEMPTS; attempt += 1) {
try {
await pool.query("SELECT 1");
await createSchema(pool);
status.connected = true;
status.lastError = null;
return;
} catch (error) {
status.connected = false;
status.lastError = String(error.message || error);
if (attempt === DB_CONNECT_ATTEMPTS) {
throw error; // 重试耗尽,向上抛出让进程退出
}
await delay(DB_CONNECT_RETRY_MS);
}
}
}
30 次 × 2 秒 = 最多等 60 秒,覆盖 MySQL 8.4 冷启动初始化的窗口。重试期间 status.connected 为 false,/api/status 会暴露这个状态,运维侧能看到 API 正在等库。重试耗尽则直接抛错让进程退出,由容器的 restart policy 兜底——比把一个连不上库的 API 留着强。
连接池本身用 mysql2/promise 的 createPool,waitForConnections: true 让瞬时并发超限时排队而不是报错,timezone: "Z" 强制按 UTC 存取时间,避免容器时区不一致把 received_at 写偏。
启动时自动建表
createSchema 在每次启动时跑一遍 CREATE TABLE IF NOT EXISTS,覆盖 measurements、alarms、users、command_acks 四张表。这样首次部署只要把环境变量配好,API 起来表就在,不用手动导 SQL。
CREATE TABLE IF NOT EXISTS measurements (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY,
device_id VARCHAR(64) NOT NULL,
topic VARCHAR(255) NOT NULL,
payload_ts BIGINT NULL, -- 设备上报的 uptime(秒)
seq BIGINT NULL, -- 设备侧自增序号
flow DOUBLE NULL,
total_value DOUBLE NULL,
velocity DOUBLE NULL,
pressure DOUBLE NULL,
temperature_json JSON NULL,
valid_mask INT NULL,
payload JSON NOT NULL, -- 原始帧整包存一份,前端可追溯
received_at DATETIME(3) NOT NULL,-- 服务器接收时间,毫秒精度
INDEX idx_measurements_device_received (device_id, received_at),
INDEX idx_measurements_device_payload_ts (device_id, payload_ts)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
两个细节:
id用BIGINT UNSIGNED,自增主键同时也是游标分页的游标,后面会用到。received_at用DATETIME(3)保留毫秒,和payload_ts(设备 uptime)分开存。两套时间戳的查询语义不同,下面单独讲。
建表之后还有一个 ensureHistoricalIdempotency 步骤,负责补唯一索引——这是历史幂等性的核心。
历史幂等性:三元唯一索引
旧方案的问题
早期版本按 (device_id, seq) 建唯一索引去重,云端确认与历史去重 那篇记过这个方案。它在“DR154 补传同一帧”时工作得很好,但碰到 MCU 重启 就崩了:
STM32 每次重启 seq 从 1 重新开始。假设重启前已经写过 (FM001, seq=1),重启后第一帧又是 (FM001, seq=1)——唯一索引拦截,ON DUPLICATE KEY UPDATE 把旧值覆盖成新值,重启前的第一帧就这么丢了。重启次数越多,历史曲线空洞越明显。
三元索引:加一个 payload_ts
seq 会归零,但 payload_ts(设备 uptime)也会归零,单加它没用。关键是 两者归零的时机不同步:seq 是每帧自增,payload_ts 是开机后持续累加。把三者拼成 (device_id, seq, payload_ts),只要两次重启后不是在完全相同的 uptime 处发出相同 seq,就不会撞键:
async function ensureHistoricalIdempotency(pool) {
// 补传可能重复到达;seq 在 MCU 重启后会从 1 重新开始,
// 因此不能只按 (device_id, seq) 去重,否则重启后的新数据会入库失败。
await ensureUniqueDeviceSeqTimestamp(
pool, "measurements",
"uniq_measurements_device_seq_payload_ts", // 新三元索引
"uniq_measurements_device_seq", // 旧二元索引,需先删
);
await ensureUniqueDeviceSeqTimestamp(
pool, "alarms",
"uniq_alarms_device_seq_payload_ts",
"uniq_alarms_device_seq",
);
}
ensureUniqueDeviceSeqTimestamp 做三件事:先清理重复行(下一节),再删旧的二元索引,最后建三元唯一索引。建索引前先查 INFORMATION_SCHEMA.STATISTICS 避免重复建:
async function ensureUniqueDeviceSeqTimestamp(pool, table, indexName, legacyIndexName) {
await removeDuplicateSeqTimestampRows(pool, table);
await dropIndexIfExists(pool, table, legacyIndexName);
const [rows] = await pool.query(
`SELECT COUNT(*) AS count
FROM INFORMATION_SCHEMA.STATISTICS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = ?
AND INDEX_NAME = ?`,
[table, indexName],
);
if (Number(rows[0]?.count) > 0) return; // 已存在则跳过
await pool.query(
`ALTER TABLE ${tableName} ADD UNIQUE KEY ${index} (device_id, seq, payload_ts)`,
);
}
写入侧:ON DUPLICATE KEY UPDATE
索引建好后,insertMeasurement 用 INSERT ... ON DUPLICATE KEY UPDATE 兜补传——同一条 (device_id, seq, payload_ts) 重复到达时不会报 1062,而是覆盖非键字段:
INSERT INTO measurements
(device_id, topic, payload_ts, seq, flow, total_value, velocity, pressure,
temperature_json, valid_mask, payload, received_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON DUPLICATE KEY UPDATE
topic = VALUES(topic),
flow = VALUES(flow),
total_value = VALUES(total_value),
velocity = VALUES(velocity),
pressure = VALUES(pressure),
temperature_json = VALUES(temperature_json),
valid_mask = VALUES(valid_mask),
payload = VALUES(payload),
received_at = VALUES(received_at)
注意 ON DUPLICATE KEY UPDATE 里没有更新 device_id/seq/payload_ts——这三个是唯一键的组成部分,改它们没意义;其余字段取 VALUES(列名) 用新值覆盖旧值。这样补传到达时最后一次的内容生效,历史曲线对应位置只剩一个点。
去重清理:removeDuplicateSeqTimestampRows
从二元索引升到三元索引之前,历史表里可能已经有重复行。直接加 UNIQUE KEY 会因为存在重复值而失败,所以必须先清理:
async function removeDuplicateSeqTimestampRows(pool, table) {
const tableName = quoteHistoricalTable(table);
await pool.query(`
DELETE dup
FROM ${tableName} dup
JOIN ${tableName} keep
ON dup.device_id = keep.device_id
AND dup.seq = keep.seq
AND dup.payload_ts = keep.payload_ts
AND dup.id > keep.id
WHERE dup.seq IS NOT NULL
AND dup.payload_ts IS NOT NULL
`);
}
这条 DELETE ... JOIN 的语义是:对每一组 (device_id, seq, payload_ts) 相同的行,保留 id 最小的那条(keep),删掉所有 id 更大的(dup)。WHERE 里的 IS NOT NULL 把 seq 或 payload_ts 为空的行排除——这些行不参与幂等约束,保留下来不碍事。
表名和索引名通过 quoteIdentifier 反引号转义,防止保留字和注入:
function quoteIdentifier(value) {
return `\`${String(value).replaceAll("`", "``")}\``;
}
replaceAll("”, ”““) 把内嵌的反引号翻倍转义,这是 MySQL 标识符转义的规范做法。这里表名是硬编码白名单(quoteHistoricalTable只允许measurements/alarms`),但转义逻辑仍按通用方式写,避免后续扩展时踩坑。
received_at vs payload_ts:两种时间戳
measurements 表里有两个时间字段,查询时必须分清语义:
received_at:API 收到这一帧的服务器时间,DATETIME(3),UTC。payload_ts:设备上报的ts字段,BIGINT。
关键点:STM32 上报的 ts 不是 Unix 时间戳,而是开机后的秒数(uptime)。设备重启后 ts 从 0 重新累加,同一个 ts=120 在两次重启之间指的是完全不同的物理时刻。如果用 payload_ts 做历史过滤,重启前的数据就再也查不到了——客户端传 from=120 会把重启后 uptime>120 的全部数据捞出来,跟用户想要的“某时刻之后”完全对不上。
所以 buildQuery 里 from/to 绑定到 received_at,并用 FROM_UNIXTIME(?) 把客户端传来的 Unix 秒转成 MySQL 时间:
function buildQuery(query) {
const clauses = [];
const params = [];
if (query.device) {
clauses.push("device_id = ?");
params.push(query.device);
}
// 使用服务器接收时间 received_at 过滤,而不是设备上报的 payload_ts。
// STM32 设备上报的 ts 是开机后的秒数(uptime),不是 Unix 时间戳,
// 设备重启后 ts 会重置,导致无法查询到重启前的历史数据。
if (query.from !== null) {
clauses.push("received_at >= FROM_UNIXTIME(?)");
params.push(query.from);
}
if (query.to !== null) {
clauses.push("received_at <= FROM_UNIXTIME(?)");
params.push(query.to);
}
// ...
}
payload_ts 仍然保留在表里和 payload JSON 里,供前端做单次会话内的相对时间展示或排序校验,但不参与跨重启的时间范围查询。建表时 (device_id, payload_ts) 索引保留,是为了将来如果固件切换到真 Unix 时间戳,可以直接复用这条索引。
游标分页:before_id
历史接口要支持“往前翻更多”,最直观的做法是 LIMIT offset, n,但深度分页时 MySQL 要扫过前 offset 行再丢弃,越往后越慢。buildQuery 用 before_id 做游标分页,ORDER BY id DESC 配合 id < ?:
// 游标分页:客户端取上一页最后一条的 id 作为 before_id 传回,按 id DESC 继续翻页。
if (query.beforeId !== null) {
clauses.push("id < ?");
params.push(query.beforeId);
}
id 是自增主键,id < ? 走主键索引,翻到第 N 页的代价和第一页一样。客户端流程:
- 第一页请求
/api/history?limit=100,拿到 100 条,记录最后一条的id。 - 第二页请求
/api/history?limit=100&before_id=<上一页最后的id>。 - 返回空列表表示到底了。
parseBigIntCursor 负责解析游标,非法值统一回 null(等价于不传),不会把 NaN 带进 SQL:
function parseBigIntCursor(value) {
if (value === null || value === undefined || value === "") return null;
const parsed = Number(value);
if (!Number.isFinite(parsed) || parsed <= 0) return null;
return Math.trunc(parsed);
}
LIMIT 内联与 WHERE 参数绑定
buildQuery 返回的 limit 是经过 parseLimit 校验并钳到 [1, MAX_LIMIT] 区间的纯数字,直接内联进 SQL;WHERE 条件里的值则全部走 ? 参数绑定:
return {
// LIMIT 只能内联数字,where 条件仍走参数绑定避免注入。
where: clauses.length ? `WHERE ${clauses.join(" AND ")}` : "",
params,
limit: query.limit,
};
function parseLimit(value) {
const parsed = Number.parseInt(value ?? "", 10);
if (!Number.isFinite(parsed) || parsed <= 0) return 100;
return Math.min(parsed, MAX_LIMIT); // MAX_LIMIT = 20000
}
为什么 LIMIT 不也用 ? 绑定?mysql2 的预处理语句里 LIMIT ? 在某些版本下会被当成字符串导致报错,且 limit 经过 parseInt + Math.min 后一定是有限正整数,内联没有任何注入风险。WHERE 里的 device_id、时间戳、before_id 都是用户可控字符串/数字,必须走参数绑定——mysql2 的 pool.query(sql, params) 会自动转义。
小结
- 连接池有限次重试:30 次 × 2 秒覆盖 MySQL 冷启动窗口,重试耗尽直接退出交给容器编排重启,不在半死状态挂着。
- 启动自动建表:
CREATE TABLE IF NOT EXISTS让首次部署零 SQL 操作,id用BIGINT UNSIGNED为游标分页铺路。 - 三元唯一索引:
(device_id, seq, payload_ts)解决 MCU 重启 seq 归零导致的去重冲突,补传走ON DUPLICATE KEY UPDATE覆盖非键字段。 - 去重清理前置:升索引前用
DELETE ... JOIN按(device_id, seq, payload_ts)收敛重复行,保留id最小的,否则加唯一索引会失败。 - 两种时间戳语义:
received_at(服务器 UTC)用于范围查询,payload_ts(设备 uptime)不跨重启使用,避免重启后历史查不到。 - 游标分页:
before_id+ORDER BY id DESC走主键索引,深度分页性能恒定。 - 注入防护:
LIMIT内联经parseInt钳制后的纯数字,WHERE全走?参数绑定,标识符走反引号转义。
后续阅读
- 云端确认与历史去重:cloud_ack 下行链路 + UNIQUE(device_id, seq)——本文三元索引的前身,记录了旧的二元去重方案及其在 MCU 重启场景下的局限。
- Docker 容器化部署与 Node.js 后端——MySQL 容器的小内存配置与 API 容器的依赖启动顺序,是本文重试逻辑的背景。
- 上行 JSON 协议与手写解析器——
tele/alarm/ack帧格式定义,ts与seq字段在设备侧的来源。