Termux上跑10层SQLite Agent Mesh:热节流感知调度实践

Termux上跑10层SQLite Agent Mesh:热节流感知调度实践
在实际开发中把多级 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输出渲染把结构化结果转成文本或 HTMLT8归档落库写入长期存档表精简临时字段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 摄氏度的值temperature430表示 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_modeWAL; PRAGMA synchronousNORMAL; PRAGMA busy_timeout5000; 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_modeWAL必须在连接初始化时执行不过多次执行也不会报错它会返回当前模式。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_modeWAL) conn.execute(PRAGMA synchronousNORMAL) conn.execute(PRAGMA busy_timeout5000) 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_modeWAL是 Mesh 能跑起来的前提没有它多个进程读同一个库时会频繁遇到锁。synchronousNORMAL在 WAL 模式下已经能保证进程崩溃后数据库不损坏性能比FULL好很多。busy_timeout5000让写操作在遇到锁时最多等待 5 秒而不是立刻抛异常。3.4 消息状态机与可能状态一个任务的完整状态流pending - running - pending进入下一层 pending - running - done最后一层完成 pending - running - pending失败重试retry_count1 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, encodingascii) 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, encodingascii) as f: freqs.append(int(f.read().strip())) except (OSError, ValueError): continue return freqs class ThermalAwareLoop: def __init__(self, low_threshold45, high_threshold60, idle_interval1.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 statuspending ORDER BY id LIMIT 1, (tier,), ).fetchone() if row is None: conn.execute(COMMIT) return None conn.execute( UPDATE tasks SET statusrunning, worker?, updated_atdatetime(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, resultNone, errorNone): 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_atdatetime(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?, workerNULL, updated_atdatetime(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_asciiFalse) if tier 1: return json.dumps({length: len(payload[content])}, ensure_asciiFalse) if tier 2: return json.dumps({enriched: True}, ensure_asciiFalse) if tier 3: return json.dumps({dedup: True}, ensure_asciiFalse) if tier 4: time.sleep(0.05) return json.dumps({computed: task_id}, ensure_asciiFalse) if tier 5: return json.dumps({checked: True}, ensure_asciiFalse) if tier 6: return json.dumps({aggregated: True}, ensure_asciiFalse) if tier 7: return json.dumps({rendered: pok/p}, ensure_asciiFalse) if tier 8: return json.dumps({archived: True}, ensure_asciiFalse) if tier 9: return json.dumps({notified: True}, ensure_asciiFalse) raise ValueError(unknown tier %d % tier) def main(): parser argparse.ArgumentParser() parser.add_argument(--tier, typeint, requiredTrue) parser.add_argument(--low, typefloat, default45.0) parser.add_argument(--high, typefloat, default60.0) args parser.parse_args() worker tier%d-%d % (args.tier, os.getpid()) loop ThermalAwareLoop(low_thresholdargs.low, high_thresholdargs.high) conn db.get_conn() print(start worker %s tier%d % (worker, args.tier), flushTrue) 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, errorstr(e)) print(task%s error%s % (task_id, e), flushTrue) else: finish_task(conn, task_id, args.tier, resultresult) print(task%s tier%d done % (task_id, args.tier), flushTrue) 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_asciiFalse) conn.execute(INSERT INTO tasks(payload) VALUES(?), (payload,)) conn.commit() print(inserted 100 tasks)新建run_mesh.py一次性拉起 10 层 Workerimport 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, flushTrue) if alive 0: break except KeyboardInterrupt: for p in procs: p.terminate()初始化数据库并启动sqlite3 mesh.db schema.sql python seed.py python run_mesh.py想单独测某一层可以只启动对应 tierpython 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.pyimport 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

最新新闻

日新闻

周新闻

月新闻