返回发布
分布式架构MQTTMySQLPython/PHP

灵动感知项目后端架构拆解

双服务进程、影子缓存、IMU 算法引擎与纯原生 MQTT 实现

双服务+PHP+MySQL全链路架构图

一、架构总览

这个项目的后端需求拆开来看有三个方向:硬件端最快 3 秒一条的上报频率要求毫秒级的数据处理能力,App 端需要稳定的 RESTful 查询接口读取健康趋势和轨迹回放,智能家居自动化场景需要低延迟的条件判定和 MQTT 下发。三个方向对延迟、吞吐量和可靠性的要求完全不同。最终落地的是两个分离的 Python 长驻进程 + 一套 PHP REST 接口层 + MySQL,外加一个 Mosquitto MQTT Broker 做消息中枢。

MQTT Broker 承载了所有实时通信。四个主题构成完整的双向通道:omega/1 接收硬件上报的传感器数据,app/fence/update 接收 App 端弹围栏修改指令,app/sync/request 接收 App 端主动全量同步请求,testtopic/1 是 Omega 下发全量状态给 App 端的通道。Alpha 通过 backend/+/config_update 接收 PHP 端发起的配置变更推送。双进程的拆分边界很清晰:Omega 约 600 行面向数据上行(设备→数据库→App),Alpha 约 180 行面向配置下行(App→数据库→设备)。Omega 的主要逻辑是接收 IMU + GPS + 温度数据后顺序执行运动分类、计步推算、围栏判定、自动化评估和小时级聚合,最后把全量状态推给 App。Alpha 只做配置同步,它的核心是一段 3 秒周期的数据库轮询加上一个内存影子缓存。

二、Omega服务:数据上行链路深度拆解

Omega 启动后订阅三个 MQTT 主题,以阻塞方式持续监听设备消息。每条消息到达后进入统一的处理入口,按主题类型分发到三个逻辑分支——设备数据上报、围栏指令更新和全量同步请求。

2.1 IMU 运动状态分类:双模式并行

硬件端存在两种上传模式。Mode 0(瞬时采样)上传原始三轴加速度和三轴陀螺仪瞬时值,由 Omega 端执行完整算法。Mode 1(平均采样)上传一个上报周期内的运动特征均值,Omega 端用不同的公式推算。两个模式的数据都以 imu 字段携带。Mode 0 的瞬时判定逻辑定义在 PetAlgorithms.analyze_motion_latest() 里:

analyze_motion_latest.py
def analyze_motion_latest(ax, ay, az, gx, gy, gz): mag = sqrt(ax² + ay² + az²) delta_g = abs(mag - 9.80665) # 重力偏差 gyro_energy = abs(gx) + abs(gy) + abs(gz) if delta_g > 4.5 or gyro_energy > 150: status = 3 # 跳跃/剧烈运动 elif delta_g > 1.2 or gyro_energy > 25: status = 1 # 行走/轻度运动 else: status = 2 # 静止/睡眠 return status, delta_g, gyro_energy

阈值 4.5 m/s² 的重力偏差和 150 deg/s 的陀螺仪能量是两个维度——delta_g 捕捉垂直方向的冲击(跳跃时的瞬间加速度变化),gyro_energy 捕捉旋转角速度的累积量(猛烈甩头或翻滚)。两个条件满足任一就判定为剧烈运动。Mode 1 的平均采样模式用 max_delta_g 和 max_gyro 做相同的三级分类,但输入的是半个周期内的峰值而非瞬时值。

2.2 计步算法:双通道互验

步数输出取 IMU 推算值和 GPS 距离推算值两者中的较大值。IMU 侧的推算依赖于运动状态分类的结果:

IMU 计步推算(Mode 1 平均采样)
# Mode 1: 活跃占比反向推算 dynamic_steps_per_min = (undiluted_gyro * 0.5) + (undiluted_delta_g * 10.0) dynamic_steps_per_min = max(15.0, min(180.0, dynamic_steps_per_min)) imu_steps = int(delta_move_min * dynamic_steps_per_min) # Mode 0: 瞬时采样 dynamic_steps_per_min = (gyro_energy * 0.5) + (delta_g * 10.0) dynamic_steps_per_min = max(15.0, min(180.0, dynamic_steps_per_min)) imu_steps = int(delta_move_min * dynamic_steps_per_min)

系数 0.5 和 10.0 是实验调试出来的经验值,分别对应陀螺仪能量和重力偏差对步频的贡献权重。陀螺仪能量在走路时大约 30-80,delta_g 大约 1.2-4.0,计算出的动态步频落在 25-100 步/分钟之间,与实际观测值基本吻合。上下饱和钳制在 15(极慢踱步)和 180(极限冲刺)之间防止异常值。GPS 侧用 Haversine 计算两次上报之间的位移距离,除以猫步幅系数 CAT_STRIDE = 0.25m 得出 GPS 推算步数。两个独立传感通道互相对照——如果猫在狭窄区域内频繁活动,GPS 位移很小但 IMU 能识别到运动,IMU 步数高;如果猫在开阔地带快速跑动,GPS 算出的步数会更准。

GPS 侧有一个两层漂移滤波器。第一层 GPS_DRIFT_THRESHOLD = 2.0m:当前位置与上次位置的距离小于 2 米时视为 GPS 噪音直接忽略,防止静止时步数持续自增。第二层 MAX_VALID_MOVE = 500.0m:超过 500 米的单次位移视为不可能发生的异常值,通常是定位突然跳点。两层都满足后才会触发 Haversine 距离转步数的逻辑。

2.3 电子围栏与环境传感器

围栏存储于 pet_config 表,包含圆心坐标和半径。每次收到有效的 GPS 数据(经度不为 0),Omega 立即用 Haversine 计算当前位置到围栏中心的距离,大于半径则 is_inside = 0。围栏有两个更新入口:App 端通过 PHP 写数据库后由 Alpha 3 秒轮询触发下发,或由 App 端直接通过 MQTT 主题 app/fence/update 发送围栏变更,Omega 的 fence handler 接收后直接写数据库、计算新位置相对于新围栏的进出状态并推回给 App 端。温度数据在 Omega 端同时处理两组:temp(环境温度传感器)和 ntc_temp(接触式高精度体温传感器)。后者写入 pet_status 表参与小时级聚合和自动化判定,前者仅作为环境参考写入。湿度同样采集后直接更新。

2.4 Automation 引擎

Omega 处理完数据后在 evaluate_automations() 中扫描 automations 表中所有 is_active = 1 的规则。每行规则包含 trigger_type(ntc_temp / temp / humidity)、operator( > / < )、threshold、action_device_id 和 action_state(ON / OFF)。Omega 用当前收到的传感器值与每条规则的阈值做比较,满足条件后执行一次去重判定——查询 smart_home_devices 表确认该设备的当前状态是否已经等于目标状态,如果不一样才触发动作。MQTT 下发的目标端口是 1884(MQTT_HA_PORT),与设备数据主端口 1883 分离,避免自动化指令与设备上报争抢 I/O。下发后在数据库写入 last_triggered 时间戳防止同一个触发源在短时间内重复触发。

2.5 小时级聚合与日级健康摘要

每次数据上报处理完后会调用 update_hourly_data()。这个方法内部维护两张小时级表。第一张 pet_hourly_stats 用累加式的 ON DUPLICATE KEY UPDATE 语法逐条增加步数、时长和温度总和/计数,计算均值。第二张 health_hourly_metrics 从 pet_hourly_stats 的当前小时行中同步 avg_temp、steps 和 active_ratio(计算方式:move_duration / 60 * 100,即运动分钟数占一小时的百分比)。两张表的分工是:前者作为累加器承接高频单次写入,后者作为查询视图面向 App 端健康页的 24 小时曲线。每小时结束时还会自动调用 update_daily_summary() 汇总当日全部 24 小时均值写入 health_daily_summary 表。日健康评分的计算规则:步数低于 2000 扣 15 分、日均体温高于 39.0°C 扣 20 分、低于 37.5°C 扣 5 分。低于 80 分标记为"需注意",高于 90 分且步数超 3500 标记为"活跃"。

三、Alpha服务:配置下行与影子缓存

Alpha 的职责很单一:保证手机端修改的设备配置安全送达到硬件端。它的实现核心是一个内存哈希表加一段 3 秒的轮询循环。

3.1 三个消息处理器

Alpha 订阅三个 MQTT 主题,对应三种场景。场景一是设备端通过拨轮或按键切换上报间隔和模式后,通过 omega/+/status 上报,Alpha 将其同步至数据库 devices 表并同时更新影子缓存。场景二是设备上线后通过 omega/+/online 通知,Alpha 从数据库拉取最新配置通过 alpha/1/cmd 下发,payload 包含间隔(秒)和模式两个字段。场景三是 PHP 端修改了配置后向 backend/+/config_update 发布消息触发 Alpha 重新下发。每个场景结束时都会更新或校验影子缓存。

3.2 影子缓存与 3 秒 Watchdog

影子缓存是一个全局 Python 字典 shadow_cache = {},key 是 device_id,value 是一组 interval + mode 的快照。设计的核心约束是:Alpha 不能把自己下发的配置再当数据库变更识别并重新下发,否则会形成无限循环。解法是影子缓存。Alpha 启动后进入一个无限循环,每 3 秒执行一次 SELECT device_id, upload_interval, upload_mode FROM devices 全量扫描。对数据库返回的每一行,检查其值是否与影子缓存中的快照一致。如果一致则跳过;如果不一致则更新影子缓存并下发指令。设备自己上报导致的数据库变更在同步后也会更新影子缓存,因此不会被轮询再次下发。

影子缓存核心逻辑(Alpha.py 主循环)
while True: conn = get_db_connection() cursor = conn.cursor(dictionary=True) cursor.execute("SELECT device_id, upload_interval, upload_mode FROM devices") for row in cursor.fetchall(): did = row['device_id'] db_interval, db_mode = row['upload_interval'], row['upload_mode'] if did not in shadow_cache: shadow_cache[did] = {"interval": db_interval, "mode": db_mode} else: cached = shadow_cache[did] if cached['interval'] != db_interval or cached['mode'] != db_mode: shadow_cache[did] = {"interval": db_interval, "mode": db_mode} # 下发给设备 client.publish("alpha/1/cmd", json.dumps({"interval": int(db_interval/1000), "mode": db_mode}), qos=1) time.sleep(3)

这个方案同时覆盖了人工介入的场景——运维人员直接登录数据库修改 devices 表的值,不做任何代码改动、不触发任何 MQTT 消息,3 秒内就会被 Alpha 侦测到并推送到设备。

四、PHP 接口层:只读不写与健康异常检测

PHP 层 15 个接口文件运行在 PHP-FPM 下。所有接口共用 db.php 建立一个 mysqli 连接,charset 设置为 utf8mb4。接口按功能分为五组:认证(login/register/send_code/reset_password)、设备绑定管理(bind/unbind/get_devices/set_active_device)、配置读写(get/set_device_config)、健康与轨迹(health_api/get_track)、智能家居与自动化(ha_api/automation_api)、社交(api_moments)。所有接口的返回格式统一为 JSON 三字段:{ code, msg, data }。

4.1 健康异常检测算法

health_api.php 是 PHP 层中唯一包含逻辑的接口,其他接口基本是参数校验后直接SQL查询输出。GET?action=get_daily_list&device_id=xxx 返回 30 日的每日摘要和趋势分析。算法步骤:遍历 30 天数据,对每一天计算其步数相对于 30 日均值乘以 0.7 的阈值,如果连续 2 天低于此阈值则触发活动量骤降警告。对体温做 3 日连续检测:如果最近 3 天的体温逐日上升(day0 > day1 > day2),则触发风险预测。同时检查当日的发热小时数(体温超过 39.1°C 的小时数),≥2 小时即判定为持续高温异常。三个检测互不排斥。PHP 端不做任何复杂的窗口计算——它每次请求时实时聚合,不以额外存储为代价。

4.2 智能家居控制流

ha_api.php 处理设备增删改。控制指令(set 动作)走的是 PHP 服务器本地执行 mosquitto_pub 系统命令向 MQTT 主机的 1884 端口发送 JSON 载荷。PHP 端先发 MQTT 再写数据库和响应 App——这是一个 fire-and-forget 模式,MQTT 下发失败不会阻塞 App 端操作。

五、数据库设计:15 张表的三层分界

petmonitor 数据库包含 15 张表,按使用模式分为三层。实时状态层(pet_status、pet_config)由 Omega 在每收到一条设备消息时高频写入,行锁开销小——每张表只有一条或数条记录,主键是 device_id 驱动的唯一行。持久化聚合层(pet_data_history、pet_hourly_stats、health_hourly_metrics、health_daily_summary)用 append-only 的周期主子表结构设计。业务管理层(users、devices、user_devices、smart_home_devices、automations、moments、moment_comments、moment_likes、verification_codes)由 PHP 接口层通过事务写入,低频、大字段、关联查询密集。三层之间不会发生锁竞争——Omega 的写入集中在实时状态表和小时聚合表上,PHP 的查询几乎只访问聚合表和管理表。

层级表名写入频率写入源
实时状态pet_status, pet_config3 秒/次Omega MQTT handler
聚合持久pet_data_history, pet_hourly_stats, health_hourly_metrics, health_daily_summary秒级聚合写入Omega update_hourly/update_daily
业务管理users, devices, user_devices, smart_home_devices, automations, moments, verification_codes用户触发PHP API

小时聚合表 pet_hourly_stats 采用复合主键 (device_id, record_hour) 搭配 ON DUPLICATE KEY UPDATE 累加语法。同一小时内接收到的多次上报数据——步数累加、体温累计加总后除以计数。温度均值字段 avg_temp 在每次更新时实时计算 temp_sum / temp_count。定位历史表 pet_data_history 使用 BIGINT UNSIGNED AUTO_INCREMENT 主键避免 32 位 INT 溢出——6 吨挖掘机和宠物定位是两个不同的量级,但历史轨迹按 3 秒一条的频次数月的积累量足以撑满 INT。

六、鸿蒙 App:纯原生 TCP Socket 实现的 MQTT 3.1.1

App 端的 MQTT 通信没有依赖任何第三方协议库。它在鸿蒙系统的原生网络套接字接口上,从零构造了 MQTT 3.1.1 的数据包编码与发送——不引入解析器,直接在 TCP 字节流中拼接协议头和数据载荷。

6.1 字节级报文构建

MQTT CONNECT 报文的组装在 executeConnect() 中完成。固定头 0x10(CONNECT 类型)+ 剩余长度后跟上协议名四个字节 0x4D 0x51 0x54 0x54("MQTT")、协议级别 0x04(MQTT 3.1.1)、连接标志 0x02(Clean Session)、Keep Alive 30 秒(0x00 0x3C)、最后是两字节长度前缀 + Client ID 的 UTF8 字节数组。剩余长度用 encodeRemainingLength() 方法按 MQTT 标准的定长编码算法生成:每字节低 7 位是有效数据,高 1 位为 1 表示继续,为 0 指示结束。

CONNECT 报文构建(GlobalContext.ets)
const mqtt3Connect = new Uint8Array([ 0x10, remainLen, 0x00, 0x04, // 协议名长度 4 0x4D, 0x51, 0x54, 0x54, // "MQTT" 0x04, // 协议级别 3.1.1 0x02, // Clean Session 0x00, 0x3C, // Keep Alive 30s (idLen >> 8) & 0xFF, idLen & 0xFF, // Client ID 长度 ...Array.from(idEncoder) // Client ID 字节 ]);

SUBSCRIBE 报文用固定头 0x82,紧接报文标识符 0x00 0x01、主题过滤器字符串长度前缀 + 字节、以及 QoS 级别字节。PUBLISH 报文用 0x30 固定头,剩余长度之后是主题长度前缀 + 主题字节 + 消息体字节。PINGREQ 是最精简的——两个字节 0xC0 0x00。

6.2 JSON 提取与状态分发

MQTT PUBLISH 报文到达时,App 端在 on('message') 回调中收到原始的字节缓冲区。处理方式是:用 TextDecoder 将 bytes 解码为 UTF-8 字符串,用 rawData.indexOf('{') 和 rawData.lastIndexOf('}') 找到 JSON 对象的起止边界提取纯 JSON 字符串跳过 MQTT 协议头。解析后的 JSON 进入 processRawData(PetJsonPayload) 方法逐一映射到 AppStorage.setOrCreate 的对应字段。温度数据的分发有明确的解耦设计:ntcTemp(接触式体温)触发通知栏推送的体温警报,而 envTemp(环境温度)仅用于 UI 显示。数据推送频率由 lastSyncTimestamp 和 5 秒 的防抖控制,在首次连接后和后续重连后只会发起一次同步请求,不会重复拉取。

6.3 断线重连状态机

状态管理通过三个布尔变量实现。isConnecting 防止并发连接尝试,isConnected 标记连接成功后禁止重复初始化。reconnectTimer 在断线时以 3000ms 延迟触发 connectToServer(isForce),其中 isForce 参数控制是否清理旧连接所有事件监听器和定时器后重新创建 TCPSocket 实例。连接成功后启动一个 30 秒周期的 setInterval 定期发送 0xC0 0x00 PINGREQ 维持连接。登出或切换账号时调用 disconnect() 清理所有资源——移除事件监听、关闭 socket、清空定时器、将连接标记重置为 false。

6.4 WGS-84 到 GCJ-02 坐标转换

CoordTransform.ts 实现了硬件 GNSS 上报的 WGS-84 国际坐标系到国内高德地图 GCJ-02 火星坐标系的非线性偏转算法。核心函数 wgs84ToGcj02(lng, lat) 首先检查经纬度是否在中国境内(经度 73.66-135.05,纬度 3.86-53.55),境外直接原值返回。境内坐标通过 transformLat 和 transformLng 两个函数计算以经度和纬度相对于中国中心点(105°E, 35°N)的偏移量为输入的非线性偏置。偏置值经过正弦级数展开模拟火星坐标系的非线性畸变——偏置量在经度方向上叠加了 20° 和 2° 周期、纬度方向上叠加了 1° 和 1/3° 周期的正弦波动以拟合真实的坐标系差异。偏置量最终乘上基于纬度的缩放因子 dLat = (dLat * 180.0) / ((A * (1 - EE)) / (magic * sqrtMagic) * PI) 后叠加到原始坐标上输出。

七、架构的弹性:五个独立域,各自降级

拆成五个独立组件的一个副产品是故障隔离。如果 Omega 进程意外退出,设备依然在持续发送数据,消息队列会暂时缓存未处理的消息,Omega 重启后按顺序恢复——在此期间手机端界面停留在最后一次推送的状态,用户察觉不到后台在恢复。如果 Alpha 进程退出,设备继续按照当前的本地配置运行,手机端暂时无法修改配置,但设备状态和健康数据的正常展示完全不受影响。如果 PHP 接口层不可用,手机端无法从服务器拉取新的历史数据,但内存中保留的最后一份数据缓存足够支撑界面正常显示。消息队列、数据库和手机端各自拥有独立的故障模式,一个组件的崩溃不会形成全链路雪崩。

这个架构中没有引入 Redis 或任何外部缓存层。这是一个有意识的设计决策——不是因为 Redis 不好,而是在这个场景中它带来的收益不构成引入复杂度的必要性。手机端拉取 30 天健康数据列表的单次查询耗时在个位数毫秒级别,已经在 MySQL 的索引覆盖范围内。有些时候,不引入一个组件是一个比引入它更好的架构决策。

技术栈汇总

上行通道 Omega

Python / paho-mqtt / IMU 运动分类 / 双通道计步 / Haversine 围栏 / 自动化条件引擎 / 3 秒写入频率

下行通道 Alpha

Python / paho-mqtt / dict 影子缓存 / 3s 全量轮询 / 死循环防护 / 设备上线即触发配置下发

接口层 PHP

15 个 REST 文件 / 仅查询无写入 / 内置 3 日温升 + 2 日骤降检测 / token 认证 / gzip 压缩

鸿蒙 App ArkTS

原生 TCP Socket 实现 MQTT 3.1.1 / 字节级 CONNECT/SUBSCRIBE/PUBLISH 组装 / AppStorage 驱动 UI

MySQL petmonitor

15 表 / 实时状态/聚合/管理三层分离 / ON DUPLICATE KEY 累加器 / BIGINT UNSIGNED 防溢出

Mosquitto

1883 设备数据 / 1884 智能家居 / 4 主线主题 / QoS 1 保证至少一次投递