6.2 KiB
数据作业与 Outbox 技术架构
本文记录 Planet v1 的任务队列、数据库 outbox 和 Earth 刷新技术边界。每类数据产品的业务用途和端到端链路见 业务架构与数据流转。
架构边界
- PostgreSQL 是 v1 的可靠任务账本和 outbox,不引入 Kafka、Celery 或 RQ。
collection_tasks记录采集、删除、清缓存和非数据库触发的 Earth refresh 任务。earth_data_change_events记录事实表或派生表变化,是 Earth 同步的可靠来源。LISTEN/NOTIFY只做低延迟唤醒;listener 仍会轮询未消费 outbox。- Redis 主要用于缓存、认证辅助、OTP / rate limit、临时日志和 WebSocket 辅助,不是可靠队列。
数据库变化同步
flowchart LR
Write["Fact or derived table write"] --> Trigger["PostgreSQL trigger"]
Trigger --> Outbox["earth_data_change_events"]
Trigger --> Notify["planet_earth_data_changes"]
Outbox --> Listener["earth_db_change_listener"]
Notify --> Listener
Listener --> Adapter["earth_layer_adapters"]
Adapter --> Cache["cache invalidation"]
Cache --> WS["earth_updates"]
- 采集、删除、定位或派生任务写入事实表或派生表。
- statement trigger 为
INSERT/UPDATE/DELETE写入 outbox。 - listener 被 notify 唤醒,或通过轮询发现未消费事件。
- listener 使用
earth_layer_adapters.py把table + source映射成 Earth layer、刷新策略和 cache pattern。 - listener 短窗口合并同 layer 事件,失效缓存并广播
earth_updates。 - 广播成功后标记 outbox consumed;失败时保留待重试。
DB 变化不再默认创建 earth_refresh 任务,因此不会被同 source 的长采集或删除任务阻塞。earth_refresh 只保留给手动清缓存和非 DB 变化刷新提示。
数据作业队列
collection_tasks 是统一 job ledger。worker 使用 PostgreSQL FOR UPDATE SKIP LOCKED 领取 queued 任务;同一 source 的写任务串行,不同 source 可并行。
| task_type | 作用 |
|---|---|
collect |
执行内置 datasource 采集 |
clear_data |
删除该 source 的采集数据和声明过的派生数据 |
clear_cache |
删除该 source 对应的 Earth / dashboard 缓存 |
earth_refresh |
非 DB 变化场景下失效智能星球图层缓存并广播刷新提示 |
接口只创建任务并返回 task_id。任务执行、进度、取消和终态由 worker 写回 collection_tasks,并通过 datasource_tasks channel 通知前端。
取消语义是“保留已提交批次”:点击停止后,后端把任务标记为 cancelling 并取消内存中的执行协程;已经提交的批次保留,未完成批次按采集器或清理任务的回滚逻辑处理。
Earth 同步事件模型
统一事件模型是 earth.layer.changed:
{
"event": "earth.layer.changed",
"action": "database_changed",
"source": "celestrak_tle",
"table": "collected_data",
"operation": "DELETE",
"layers": ["satellites"],
"refresh_strategy": "clear_then_reload",
"records_processed": 11125,
"occurred_at": "2026-05-25T10:20:30Z"
}
| strategy | 用途 |
|---|---|
clear_then_reload |
先清前端本地图层对象,再强制重拉接口。删除数据时优先使用。 |
reload |
保留旧对象直到新数据返回,适合定位、元数据或非破坏性更新。 |
delta |
只用于 earth_interactables,按 id upsert 或 remove。 |
接口在真实 0 数据时必须返回 200 和空集合;只有真实接口异常才返回 5xx。前端收到删除事件后,如果重拉失败,应保持已清空状态并显示轻量错误,不恢复旧对象。
Layer Adapter 约定
earth_layer_adapters.py 是 source、派生表、Earth layer、缓存和刷新策略的唯一注册表。新图层只应新增 adapter,不应在按钮 handler、采集器或前端分支里手写同步逻辑。
Adapter 必须声明:
- source 或 table 由哪个 Earth layer 消费。
- 需要清理哪些 Earth cache key pattern。
clear_data删除 source 时是否需要同时删除 owned 派生表。- 该 layer 的默认刷新策略。
删除 source 时,clear_data 作业先删除 collected_data.source = <source>,再根据 adapter 删除 owned 派生表。直接修改派生表也会触发 outbox,所以后台任务、管理接口和 SQL 修复脚本只要落到事实表或派生表,Earth 都能感知变化。
运维排障
检查 outbox 是否堆积:
SELECT id, table_name, operation, source, occurred_at
FROM earth_data_change_events
WHERE consumed_at IS NULL
ORDER BY id
LIMIT 20;
检查触发器是否存在:
SELECT tgname, tgrelid::regclass
FROM pg_trigger
WHERE tgname LIKE 'tr_planet_%_changed_%'
ORDER BY 2, 1;
常用日志事件:
earth.db_changes.connected:listener 已连接 PostgreSQL 并开始监听。earth.db_changes.outbox_polled:轮询到了未消费 outbox。earth.db_changes.broadcasted:已产生 Earth 刷新广播。data_job.started/data_job.completed:任务执行状态。
如果 Earth 没更新,按顺序检查:事实表是否变化、outbox 是否消费、adapter 是否覆盖对应 source/table、listener 是否在线、前端 WebSocket 是否连接、visualization 接口是否返回 200 空集合或新数据。
Kafka-ready 边界
业务代码不直接依赖具体队列实现,而是通过这些边界组织:
JobQueue:提交、领取、取消、完成数据作业。DataChangeBus:发布数据库事实变化。EarthLayerAdapterRegistry:声明 source、layer、缓存和派生数据关系。
需要 Kafka 的信号:
- 多个独立服务需要消费同一批数据变化。
- AIS、BGP 或传感器流达到持续高吞吐。
- 需要 consumer group、事件回放、跨服务解耦。
需要 Spark 的信号:
- 历史数据到千万或亿级,PostgreSQL 聚合开始吃力。
- 需要跨源、长时间窗口、空间时间关联分析。
- 原始数据进入 Parquet / Iceberg / Delta 等湖仓,并开始生产离线派生数据产品。
若目标是秒级连续流计算,优先评估 Flink;Spark 更适合批量或微批分析。