为什么 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 后发查询会拿到 ECONNREFUSEDinitializePool 做了有限次重试,而不是无限等待:

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.connectedfalse/api/status 会暴露这个状态,运维侧能看到 API 正在等库。重试耗尽则直接抛错让进程退出,由容器的 restart policy 兜底——比把一个连不上库的 API 留着强。

连接池本身用 mysql2/promisecreatePoolwaitForConnections: true 让瞬时并发超限时排队而不是报错,timezone: "Z" 强制按 UTC 存取时间,避免容器时区不一致把 received_at 写偏。

启动时自动建表

createSchema 在每次启动时跑一遍 CREATE TABLE IF NOT EXISTS,覆盖 measurementsalarmsuserscommand_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

两个细节:

  • idBIGINT UNSIGNED,自增主键同时也是游标分页的游标,后面会用到。
  • received_atDATETIME(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

索引建好后,insertMeasurementINSERT ... 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 NULLseqpayload_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 的全部数据捞出来,跟用户想要的“某时刻之后”完全对不上。

所以 buildQueryfrom/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 行再丢弃,越往后越慢。buildQuerybefore_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 页的代价和第一页一样。客户端流程:

  1. 第一页请求 /api/history?limit=100,拿到 100 条,记录最后一条的 id
  2. 第二页请求 /api/history?limit=100&before_id=<上一页最后的id>
  3. 返回空列表表示到底了。

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 都是用户可控字符串/数字,必须走参数绑定——mysql2pool.query(sql, params) 会自动转义。

小结

  • 连接池有限次重试:30 次 × 2 秒覆盖 MySQL 冷启动窗口,重试耗尽直接退出交给容器编排重启,不在半死状态挂着。
  • 启动自动建表CREATE TABLE IF NOT EXISTS 让首次部署零 SQL 操作,idBIGINT 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 全走 ? 参数绑定,标识符走反引号转义。

后续阅读