在实际开发中,把多级 Agent Mesh 跑在 Termux 里,和跑在云服务器上最大的差别不在 Python 语法,也不在 SQLite 的能力,而在热节流(thermal throttling)。SQLite 在移动端通常几十 GB 的数据规模下性能绰绰有余,Termux 又提供了完整的 Linux 用户态环境,所以在手机上搭建一个分层任务处理网格是可行的。问题在于:Android 内核会在 SoC 或电池温度越过阈值后主动降低 CPU 频率、限制部分核心,表现出来就是任务变慢、日志出现长停顿、高负载进程被反复挤压。这篇文章用一个 10 层 SQLite Agent Mesh 的最小实现,把“架构怎么搭、SQLite 怎么当协调层、温度怎么读、调度怎么躲开热节流”这条线完整讲一遍。
适合读者:想在手机或安卓平板上跑本地自动化任务的开发者,对 SQLite 并发、移动端 Linux、任务调度感兴趣的工程师,以及想理解热节流机制并在代码层面对抗它的实践者。示例代码用 Python 3,全部可以在 Termux 内直接运行;真实设备上你需要根据机型、热区路径和任务内容调整参数。
1. 先理解 Agent Mesh 为什么需要 SQLite 当协调层
1.1 Agent Mesh 和中心化调度的区别
Agent Mesh 是一组独立运行的 Worker 进程,它们各自负责一类任务,彼此不直接调用接口,而是通过共享状态来协作。和“一个调度中心下发任务,多个执行器回报结果”的传统架构不同,Mesh 里没有单一调度节点,任何一个 Worker 挂掉,其他 Worker 仍然可以从共享状态里继续消费任务。
这个模式放在 Termux 里有一个实际好处:手机上的进程随时可能被 Android 杀掉,如果设计成强依赖中心调度器,调度器一死整个系统就停摆;如果用 SQLite 作为协调层,每次任务的状态都落库,Worker 重启后可以继续接手未完成任务,系统韧性会好很多。
1.2 SQLite 能承担协调层吗
很多人一听“多进程 + 协调层 + 消息队列”,第一反应是上 Redis 或 RabbitMQ。但在 Termux 里,这两者的部署成本和资源占用都不低,而且以手机端的数据量来说,SQLite 完全够用。
SQLite 作为协调层的关键能力有三个:
- 事务。
BEGIN IMMEDIATE可以保证“查待办任务 -> 标记为 running”这个动作是原子的,多个 Worker 同时抢同一批任务时不会重复处理。 - WAL 模式。写入时只需要追加到
-wal文件,Reader 不会被 Writer 阻塞,非常适合“多个 Agent 同时读队列、少量 Agent 写结果”的 Mesh 场景。 - 单文件。整个任务状态、事件日志、结果数据都放在一个
.db文件里,备份和迁移非常简单,直接cp就行。
它的限制也要说清楚:SQLite 同一时刻只有一个写事务,如果 10 个层同时高频写库,会出现database is locked。所以在设计中,我们要求每个 Agent 处理任务时只做一次短事务,并且用busy_timeout让写操作等待锁释放。
1.3 为什么要组织成 10 层而不是一层
单层 Agent 把所有逻辑写在一起,优点是简单,缺点也很明显:任务一旦处理到一半,进程被 Android 杀掉,你无法知道这个任务到底进行到哪一步;而且所有逻辑都压在同一个进程里,发热和功耗集中,更容易触发热节流。
10 层结构就是把一个完整任务拆成 10 个阶段,每个阶段由一个独立 Worker 消费。这样有两个直接收益:
- 可观察性。每个任务停在哪一层、哪一台 Worker 在处理、失败时重试了几次,都可以从数据库里查出来。
- 热管理。不同层可以设置不同的运行节奏,核心计算层 T4 跑慢一点,校验和存档层 T0、T9 可以保持轻负载,避免所有进程同时抢 CPU。
下面这 10 层是示例划分,真实项目可以按业务裁剪:
| 层 | 名称 | 职责 |
|---|---|---|
| T0 | 入口校验 | 接收原始任务,检查必填字段 |
| T1 | 解析归范 | 把 JSON 拆分字段,统一时间格式 |
| T2 | 数据富化 | 查本地参考表,补齐关联信息 |
| T3 | 去重归一 | 按业务主键去重,生成批次号 |
| T4 | 核心计算 | 执行主要计算逻辑,耗时最长 |
| T5 | 结果校验 | 检查结果范围、空值、一致性 |
| T6 | 聚合汇总 | 按维度聚合,生成汇总记录 |
| T7 | 输出渲染 | 把结构化结果转成文本或 HTML |
| T8 | 归档落库 | 写入长期存档表,精简临时字段 |
| T9 | 通知清理 | 生成通知记录,删除中间数据 |
2. Termux 环境准备与依赖安装
2.1 安装 Termux 并初始化基础环境
Termux 的安装渠道建议直接使用官方说明。安装完成后打开应用,先执行更新,否则后续pkg install可能因为软件源索引过期而失败:
pkg update pkg upgrade -y接着申请存储权限,这一步会给 Termux 暴露~/storage目录,方便把数据库文件移动到手机共享目录查看:
termux-setup-storage如果希望任务在息屏后继续运行,还需要安装 Termux 的扩展工具包:
pkg install termux-tools termux-apitermux-wake-lock在termux-tools里,它可以向系统申请持锁,让 CPU 在息屏后不会立刻进入深度休眠。
注意:root 并不是必需的。热区温度在许多设备上不开放给普通用户读取,但如果读不到,本文后面会给出替代方案。
2.2 安装 Python、SQLite 与常用工具
执行下面的命令安装运行环境:
pkg install python sqlite clang libffipython用于运行 Agent 代码,sqlite提供sqlite3命令行工具,clang和libffi是为了某些 Python 原生扩展能编译安装。本项目只用标准库sqlite3,不需要额外pip包。
验证安装:
python --version sqlite3 --version输出应符合 Termux 当前软件源里的版本号。如果sqlite3命令找不到,检查pkg list-installed里是否包含 sqlite。
2.3 确认手机热区和 CPU 频率接口可读
热节流的核心数据源是 sysfs,路径通常长这样:
/sys/class/thermal/thermal_zone0/temp /sys/class/thermal/thermal_zone1/temp /sys/devices/system/cpu/cpufreq/policy0/scaling_cur_freq先手动读一下:
cat /sys/class/thermal/thermal_zone0/temp cat /sys/devices/system/cpu/cpufreq/policy0/scaling_cur_freq第一行是温度,单位通常是毫摄氏度,比如43000表示 43 摄氏度,但也有少数设备直接输出43。第二行是当前频率,单位是 kHz,比如1804800表示约 1.8 GHz。
如果这些路径都不存在,说明你的设备没有把热区暴露给普通用户。可以退而求其次,读电池温度:
dumpsys battery | grep temperature但注意dumpsys输出的是 0.1 摄氏度的值,temperature=430表示 43.0 摄氏度。生产级实现最好同时支持多种温度来源。
2.4 学习环境与生产环境的差异
在 Termux 里做实验,和真正在生产环境跑是有差距的,这个差异要在动手前就说清楚:
| 项目 | 学习环境 | 生产环境 |
|---|---|---|
| 数据量 | 几千条任务 | 可能要分段处理海量任务 |
| 运行时长 | 几分钟到几小时 | 需要长期稳定运行 |
| 崩溃恢复 | 手动重启 | 需要守护进程和开机自启 |
| 监控 | 打印日志 | 指标采集、日志轮转、告警 |
| 并发 | 每层 1 个 Worker | 每层多个 Worker,但要限制总并发 |
| 数据库备份 | 复制文件 | 定时 checkpoint 和备份策略 |
手机跑 Agent Mesh 更适合做学习、原型验证和轻量边缘任务,不建议作为高吞吐在线服务。这个定位决定了后面很多设计取舍。
3. 10 层 Agent Mesh 的最小可运行架构
3.1 用一张任务表完成所有层的协调
为了让代码足够简单,这里不采用“每层一张表”的设计,而是只维护一张tasks表,用current_tier字段表示任务当前处于第几层。
每个 Worker 只做三件事:
- 从
tasks表里找current_tier = 自己的层号且status = 'pending'的任务。 - 原子地把它改成
running,防止其他 Worker 抢到。 - 处理完成后,把结果写回
result,把current_tier加 1,重新置为pending,供下一层消费。
最后一层 T9 处理完后,把任务状态置为done。这就是整个 Mesh 的闭环。
3.2 数据库表设计
schema.sql内容如下:
PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; PRAGMA busy_timeout=5000; CREATE TABLE IF NOT EXISTS tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, payload TEXT NOT NULL, current_tier INTEGER NOT NULL DEFAULT 0, status TEXT NOT NULL DEFAULT 'pending', retry_count INTEGER NOT NULL DEFAULT 0, worker TEXT, result TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')), updated_at TEXT NOT NULL DEFAULT (datetime('now')) ); CREATE INDEX IF NOT EXISTS idx_tasks_claim ON tasks(current_tier, status, id); CREATE TABLE IF NOT EXISTS task_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, task_id INTEGER NOT NULL, tier INTEGER NOT NULL, worker TEXT, event TEXT NOT NULL, message TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')) );关键点:
current_tier和status组成联合索引,让“按层抢任务”的查询走索引,避免全表扫描。retry_count记录失败重试次数,达到上限后标记为failed。task_events是审计日志表,每个任务被谁领取、由哪层处理完、失败原因是什么,都留一条记录,方便排查热节流导致的任务积压。- 第一条
PRAGMA journal_mode=WAL必须在连接初始化时执行,不过多次执行也不会报错,它会返回当前模式。
3.3 数据库访问封装
新建db.py,统一管理连接和事件写入:
import sqlite3 DB_PATH = "mesh.db" def get_conn(): conn = sqlite3.connect(DB_PATH) conn.row_factory = sqlite3.Row conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA synchronous=NORMAL") conn.execute("PRAGMA busy_timeout=5000") return conn def log_event(conn, task_id, tier, worker, event, message=""): conn.execute( "INSERT INTO task_events(task_id, tier, worker, event, message) " "VALUES(?, ?, ?, ?, ?)", (task_id, tier, worker, event, message), )journal_mode=WAL是 Mesh 能跑起来的前提,没有它,多个进程读同一个库时会频繁遇到锁。synchronous=NORMAL在 WAL 模式下已经能保证进程崩溃后数据库不损坏,性能比FULL好很多。busy_timeout=5000让写操作在遇到锁时最多等待 5 秒,而不是立刻抛异常。
3.4 消息状态机与可能状态
一个任务的完整状态流:
pending -> running -> pending(进入下一层) pending -> running -> done(最后一层完成) pending -> running -> pending(失败重试,retry_count+1) pending -> running -> failed(重试次数达到上限)pending是“可以被领取”的状态,running是“正在被某层处理”的状态。任何进程在启动后,只需要扫描pending状态的任务,就能从上次中断的位置继续。这就是为什么 SQLite 作为协调层能天然支持崩溃恢复。
4. 核心代码实现
4.1 温度读取与热感知调度模块
新建thermal.py,它负责两件事:读取当前最高温度,以及根据温度决定 Agent 是否应该暂停。
import glob import time THERMAL_GLOBS = [ "/sys/class/thermal/thermal_zone*/temp", ] FREQ_GLOBS = [ "/sys/devices/system/cpu/cpufreq/policy*/scaling_cur_freq", ] def read_temperature_celsius(): temps = [] for path in sorted(glob.glob("/sys/class/thermal/thermal_zone*/temp")): try: with open(path, "r", encoding="ascii") as f: raw = int(f.read().strip()) except (OSError, ValueError): continue if raw > 1000: temp = raw / 1000.0 else: temp = float(raw) if temp < 0: continue temps.append(temp) return max(temps) if temps else None def read_current_freqs_khz(): freqs = [] for path in sorted(glob.glob("/sys/devices/system/cpu/cpufreq/policy*/scaling_cur_freq")): try: with open(path, "r", encoding="ascii") as f: freqs.append(int(f.read().strip())) except (OSError, ValueError): continue return freqs class ThermalAwareLoop: def __init__(self, low_threshold=45, high_threshold=60, idle_interval=1.0): self.low_threshold = low_threshold self.high_threshold = high_threshold self.idle_interval = idle_interval def wait_for_slot(self): temp = read_temperature_celsius() if temp is None: time.sleep(self.idle_interval) return if temp >= self.high_threshold: sleep_seconds = min(30.0, (temp - self.high_threshold) * 2.0) time.sleep(sleep_seconds) elif temp > self.low_threshold: time.sleep(1.0) else: time.sleep(self.idle_interval)这是整篇文章最核心的模块。low_threshold和high_threshold默认是 45 和 60 摄氏度,实际值要根据不同手机调整。超过high_threshold后,暂停时间按超出程度线性增加,最高 30 秒,让 SoC 有散热窗口;处于中间区间时每次至少等 1 秒;低于低阈值时则恢复正常轮询。
注意:温度单位在不同设备上不一致。这里用
raw > 1000判断是毫摄氏度还是摄氏度,大多数设备适用,但个别设备可能输出430表示 43.0 摄氏度,落地前先手动确认一遍你手机上的原始值。
4.2 Agent 主循环与原子抢任务
新建agent.py,它是每一层 Worker 的通用实现:
import argparse import json import os import sqlite3 import time import db from thermal import ThermalAwareLoop TOTAL_TIERS = 10 def claim_next_task(conn, tier, worker): conn.execute("BEGIN IMMEDIATE") row = conn.execute( "SELECT id FROM tasks " "WHERE current_tier=? AND status='pending' " "ORDER BY id LIMIT 1", (tier,), ).fetchone() if row is None: conn.execute("COMMIT") return None conn.execute( "UPDATE tasks SET status='running', worker=?, " "updated_at=datetime('now') WHERE id=?", (worker, row["id"]), ) db.log_event(conn, row["id"], tier, worker, "claim") conn.execute("COMMIT") return row["id"] def finish_task(conn, task_id, tier, result=None, error=None): row = conn.execute("SELECT * FROM tasks WHERE id=?", (task_id,)).fetchone() if error is not None: retry_count = row["retry_count"] + 1 status = "failed" if retry_count >= 3 else "pending" conn.execute( "UPDATE tasks SET status=?, retry_count=?, result=?, " "updated_at=datetime('now') WHERE id=?", (status, retry_count, error, task_id), ) db.log_event(conn, task_id, tier, None, "error", error) else: next_tier = tier + 1 if tier + 1 < TOTAL_TIERS else TOTAL_TIERS status = "done" if next_tier == TOTAL_TIERS else "pending" conn.execute( "UPDATE tasks SET status=?, current_tier=?, result=?, " "worker=NULL, updated_at=datetime('now') WHERE id=?", (status, next_tier, result, task_id), ) db.log_event(conn, task_id, tier, row["worker"], "finish", "next_tier=%d" % next_tier) conn.commit() def process_task(tier, task_id): conn = db.get_conn() row = conn.execute("SELECT * FROM tasks WHERE id=?", (task_id,)).fetchone() payload = json.loads(row["payload"]) if tier == 0: if "content" not in payload: raise ValueError("missing content") return json.dumps({"ok": True}, ensure_ascii=False) if tier == 1: return json.dumps({"length": len(payload["content"])}, ensure_ascii=False) if tier == 2: return json.dumps({"enriched": True}, ensure_ascii=False) if tier == 3: return json.dumps({"dedup": True}, ensure_ascii=False) if tier == 4: time.sleep(0.05) return json.dumps({"computed": task_id}, ensure_ascii=False) if tier == 5: return json.dumps({"checked": True}, ensure_ascii=False) if tier == 6: return json.dumps({"aggregated": True}, ensure_ascii=False) if tier == 7: return json.dumps({"rendered": "<p>ok</p>"}, ensure_ascii=False) if tier == 8: return json.dumps({"archived": True}, ensure_ascii=False) if tier == 9: return json.dumps({"notified": True}, ensure_ascii=False) raise ValueError("unknown tier %d" % tier) def main(): parser = argparse.ArgumentParser() parser.add_argument("--tier", type=int, required=True) parser.add_argument("--low", type=float, default=45.0) parser.add_argument("--high", type=float, default=60.0) args = parser.parse_args() worker = "tier%d-%d" % (args.tier, os.getpid()) loop = ThermalAwareLoop(low_threshold=args.low, high_threshold=args.high) conn = db.get_conn() print("start worker %s tier=%d" % (worker, args.tier), flush=True) while True: loop.wait_for_slot() task_id = claim_next_task(conn, args.tier, worker) if task_id is None: time.sleep(1.0) continue try: result = process_task(args.tier, task_id) except Exception as e: finish_task(conn, task_id, args.tier, error=str(e)) print("task=%s error=%s" % (task_id, e), flush=True) else: finish_task(conn, task_id, args.tier, result=result) print("task=%s tier=%d done" % (task_id, args.tier), flush=True) if __name__ == "__main__": main()这里有两个容易出错的地方:
claim_next_task里用BEGIN IMMEDIATE,它会立刻申请写锁,再执行SELECT和UPDATE,从而避免两个 Worker 同时查到同一条pending任务。- 处理阶段可能耗时较长,所以
process_task里重新打开一个连接读取完整任务内容。领取时只锁很短的时间,释放后其他 Worker 可以继续抢任务,这是控制锁竞争的关键。
4.3 初始化数据与启动全部层
新建seed.py,插入 100 条示例任务:
import json import db conn = db.get_conn() for i in range(100): payload = json.dumps({"content": "task-%03d" % i, "seq": i}, ensure_ascii=False) conn.execute("INSERT INTO tasks(payload) VALUES(?)", (payload,)) conn.commit() print("inserted 100 tasks")新建run_mesh.py,一次性拉起 10 层 Worker:
import subprocess import sys import time procs = [] for tier in range(10): p = subprocess.Popen([sys.executable, "agent.py", "--tier", str(tier)]) procs.append(p) try: while True: time.sleep(5) alive = sum(1 for p in procs if p.poll() is None) print("alive=%d/10" % alive, flush=True) if alive == 0: break except KeyboardInterrupt: for p in procs: p.terminate()初始化数据库并启动:
sqlite3 mesh.db < schema.sql python seed.py python run_mesh.py想单独测某一层,可以只启动对应 tier:
python agent.py --tier 4 --low 40 --high 55--low和--high是温度阈值,这个例子把核心计算层的阈值调低,表示 T4 对温度更敏感,一旦超过 55 度就放慢节奏。
5. 热节流机制与监测方法
5.1 热节流是怎么发生的
Android 设备内部有多种温度传感器,分布在 SoC 内部、电池、充电模块等位置。内核的 thermal 框架会持续采样这些传感器,当温度超过预设阈值时,会触发对应的冷却设备。
最常见的冷却手段就是cpufreq调速器主动降低 CPU 频率,严重时还会停掉部分小核或大核。这和你代码里写多少线程无关,属于系统级的强制干预。标题里提到的prochot ext,在部分 x86 和 SoC 上指外部硬件信号触发 PROCHOT 事件,通俗理解就是硬件在说“太热了,立刻降低功耗”。移动端虽然不一定都叫这个名字,但本质一致:防护机制优先于性能。
所以“不触发热节流”的目标,不是修改内核参数去屏蔽保护,而是通过调度让设备温度始终低于触发阈值。
5.2 温度与频率记录脚本
为了观察节流,我们需要把温度和频率同时记录下来。新建temp_logger.py:
import csv import time from thermal import read_temperature_celsius, read_current_freqs_khz with open("thermal_log.csv", "w", newline="") as f: writer = csv.writer(f) writer.writer