Piko - Data-Oriented Async Task Orchestrator
一个面向 数据任务(ETL / 同步 / 扫描 / 归档 / 监控) 的异步并发编排框架:
代码里写任务(Job),数据库里配调度与参数(Schedule/Config),运行时自动热更新,天然支持异步高并发与可观测性。
目录
- 为什么是 Piko
- 核心特性
- Watcher 机制与正确热更新方式
- 30 秒上手
- 基础概念速览
- 大量示例:异步 / 并发 / 异步并发
- 示例 1:最小 Job + DB 调度
- 示例 2:并发抓取 HTTP(有界并发 + 超时)
- 示例 3:并发扇出/扇入(fan-out / fan-in)
- 示例 4:生产者-消费者(asyncio.Queue 背压)
- 示例 5:异步 I/O + CPU 混合流水线(ProcessPool MapReduce)
- 示例 6:把同步阻塞库变成“异步可并发”
- 示例 7:有状态任务(水位线 + 自动补跑)
- 示例 8:资源依赖注入 Resource(连接池/客户端自动释放)
- 示例 9:自定义 Resource(你想注入什么都行)
- 示例 10:持久化写入(队列缓冲 + 批量写 + 磁盘兜底)
- 示例 11:TypedSink 类型路由(不同模型不同写法)
- 示例 12:定时任务三种触发器(cron/interval/date)
- 示例 13:动态配置热更新(灰度生效 effective_from)
- 示例 14:多实例部署(Leader Election)
- 数据库迁移
- 配置
- 运维端点与可观测性
- 项目结构建议
- FAQ
为什么是 Piko
当你在生产里做“数据任务/同步任务”时,通常会遇到这些痛点:
- 任务是 asyncio 的,但调度、并发、资源生命周期、幂等锁、补跑、观测……要自己拼很多代码
- 任务参数经常变,想做到 在线改配置、灰度生效、无需发版
- 多实例部署时,需要 只让一个实例真正执行调度(Leader/Follower)
- 高并发抓取/同步时,需要 有界并发(不炸库、不打爆下游、不 OOM)
- 任务写入落库/发 MQ 等常常成为瓶颈,需要 队列缓冲、批量写、背压、兜底
Piko 把这些能力做成“框架默认能力”,你只需要专注写业务 Job。
核心特性
Piko 的核心设计: ✅ 代码注册任务(白名单模式) + ✅ 数据库配置调度/参数 + ✅ 运行时 Reconcile 热更新
- 异步任务模型:任务必须是
async def,天然支持 asyncio 高并发 - DB 驱动调度:
scheduled_job表配置 cron/interval/date,job_config表配置参数(支持版本、灰度生效) - 运行时自动热更新:
ConfigWatcher周期性 reconcile DB → 内存缓存 → APScheduler(无需重启) - 幂等锁:
job_lock防止同一 job 在同一 scheduled_time 上重复执行 - 有状态任务:维护水位线
last_data_time,支持 自动补跑(Backfill) - 资源依赖注入(Resource):用
asynccontextmanager管理连接池/客户端,任务结束自动释放 - CPU 计算池(多进程):
CpuManager支持submit/map_reduce,绕过 GIL 做 CPU 并行 - 持久化写入引擎:
PersistenceWriter队列缓冲 + 批量聚合 + 背压 + 磁盘兜底恢复 - 内置运维 API:
/healthz/readyz/metrics(FastAPI + Prometheus) - 分布式 Leader Election:基于 DB 租约 + CAS 乐观锁,多实例下只有 Leader 执行调度
Watcher 机制与正确热更新方式
这一节专门说明 Piko 的热更新是如何工作的,以及你应该如何正确更新 job_config / scheduled_job,避免“我明明改了 DB,但任务还按旧配置跑”的问题。
1) Watcher 的本质:DB 轮询 + Reconcile(不是文件系统 watch)
Piko 的 ConfigWatcher 不是 inotify/watchdog 那类“文件变化触发”,而是一个后台协调循环:
- 周期性从 DB 读取 期望状态(
scheduled_job/job_config) - 与进程内 当前状态(APScheduler / ConfigCache)对比
- 把差异“收敛”回 DB 的状态(reconcile)
因此热更新的延迟 ≈ poll_interval_s ± jitter(轮询间隔 + 抖动)。
2) Watcher 每轮做什么(理解热更新链路)
每一轮 reconcile 主要做两件事:
-
同步 job 参数(job_config → ConfigCache)
job_config.config_json写入内存缓存(ConfigCache)- 支持
effective_from:如果配置的生效时间在未来,则暂不写入缓存(灰度/延时生效) - 任务真正执行时会从 ConfigCache 取
ctx["config"],因此 下一次触发执行就会使用新参数
-
同步 job 调度(scheduled_job → APScheduler)
- 对比 DB 中 enabled 的 job 列表与 APScheduler 当前 job 列表:增、删、改
- 重要:Piko 判断“是否需要更新触发器”的依据是
scheduled_job.version- watcher 会把
version写进 APScheduler job 的name(例如v1 / v2) - 只有当版本标签变化时才会
reschedule_job(...) - 这意味着:仅修改
schedule_expr等配置,而不 bumpversion,调度不会热更新
- watcher 会把
3) 如何正确更新 job_config(参数热更新)
3.1 立即生效(推荐写法:UPSERT + version+1)
INSERT INTO job_config(job_id, schema_version, config_json, version)
VALUES ("your_job_id", 1, '{"key":"value"}', 1)
ON DUPLICATE KEY UPDATE
config_json='{"key":"value"}',
version = version + 1,
effective_from = NULL;
说明:
version用于审计/追踪(便于定位某次运行到底用了哪个配置版本)- 参数热更新是否生效,通常不依赖
version判断(watcher 每轮会同步 config),但建议保持 version +1 的规范 effective_from = NULL表示立即生效
3.2 灰度/延时生效(effective_from)
UPDATE job_config
SET config_json='{"key":"value"}',
effective_from = DATE_ADD(UTC_TIMESTAMP(6), INTERVAL 10 MINUTE),
version = version + 1
WHERE job_id="your_job_id";
说明:
- 在
effective_from到来之前 watcher 会跳过该配置写入缓存 - 到点后下一轮 reconcile 才会真正生效(仍受
poll_interval_s影响)
4) 如何正确更新 scheduled_job(调度热更新,必须 bump version)
4.1 修改调度表达式时:一定要同时 bump version
例如把 {"minute":"*"} 改成每 4 小时整点(00 分)运行:
UPDATE scheduled_job
SET
schedule_expr = '{"hour":"*/4","minute":"0"}',
version = version + 1
WHERE job_id = "your_job_id";
或者用 UPSERT:
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("your_job_id", "cron", '{"hour":"*/4","minute":"0"}', 1, 1)
ON DUPLICATE KEY UPDATE
enabled = 1,
schedule_expr = '{"hour":"*/4","minute":"0"}',
version = version + 1;
✅ 重点结论(请直接当成规则记住):
- job_config:改参数即可生效(受 effective_from / poll_interval 影响),建议 version+1 便于审计
- scheduled_job:改 schedule_expr / schedule_type / timezone 等调度相关字段时,必须 version+1,否则不会 reschedule
4.2 禁用/删除任务
-- 禁用(推荐:禁用而不是删行,便于审计)
UPDATE scheduled_job
SET enabled = 0, version = version + 1
WHERE job_id = "your_job_id";
5) 常见排查 Checklist(“我改了但没生效”)
- 你改的是 正确的库/正确的环境 吗?(容器内 DSN、配置文件)
- 你改的是 正确的 job_id 吗?(避免把另一个仍每分钟的 job 当成目标 job)
scheduled_job调度变更时,你是否 version +1?job_config是否设置了未来的effective_from导致暂未生效?- 你是否在多实例部署?只有 Leader 负责调度(Follower standby),确认你观察的是 Leader 的日志/行为
- watcher 轮询间隔
poll_interval_s是否过大?(热更新不是实时触发)
30 秒上手
你需要一个 MySQL(Piko 用它存调度、配置、幂等锁、运行记录等元数据)。
1) 安装
# uv(推荐)
uv pip install piko-cucc
2) 配置 MySQL DSN
export PIKO_MYSQL_DSN="mysql+aiomysql://user:pass@127.0.0.1:3306/piko?charset=utf8mb4"
也可以写到
settings.toml / piko.toml,或用PIKO_SETTINGS_PATH指向自定义配置文件。
3) 初始化数据库 schema
piko db upgrade # 创建全部表并写入 schema 版本
piko db current # 确认输出 schema_v1
PikoApp 启动时会校验 schema 版本,未初始化或版本不匹配会直接启动失败。详见 数据库迁移。
4) 写一个 Job,然后跑起来
# app.py
from piko import PikoApp
app = PikoApp(name="demo")
api_app = app.api_app
@app.job(job_id="hello_job")
async def hello(ctx, scheduled_time):
print("hello", ctx["run_id"], scheduled_time)
if __name__ == "__main__":
app.run()
启动:
python app.py
# 或 ASGI:
# uvicorn app:api_app --reload
然后在 DB 插入一条调度:
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("hello_job", "interval", '{"seconds": 10}', 1, 1)
ON DUPLICATE KEY UPDATE enabled=1, schedule_expr='{"seconds":10}', version=version+1;
基础概念速览
Job(任务)
- 用
@app.job(job_id=...)注册(白名单) - 函数签名:
async def handler(ctx, scheduled_time, **resources) ctx里至少包含:ctx["run_id"]:本次执行记录 ID(job_run.run_id)ctx["job_id"]:任务 IDctx["config"]:任务配置(dict 或 Pydantic Model)- 若是有状态任务:
ctx["data_interval"](DataInterval)
调度与参数(DB)
scheduled_job:配置触发器(cron/interval/date)job_config:配置参数(版本化、灰度生效effective_from)job_run:执行记录(状态、耗时、错误)job_lock:幂等锁(同 job_id + scheduled_time 只允许一个实例执行)
示例:异步 / 并发 / 异步并发
下面的示例尽量都遵循同一个模式:你写 job,Piko 负责调度/幂等/资源/观测。你可以直接复制这些片段到自己的 jobs.py 中使用。
示例 1:最小 Job + DB 调度
from piko import PikoApp
app = PikoApp(name="mini")
api_app = app.api_app
@app.job(job_id="mini_job")
async def mini_job(ctx, scheduled_time):
# ctx["config"] 默认为 {},如果你没在 job_config 表里配置
print("run_id=", ctx["run_id"], "scheduled_time=", scheduled_time, "config=", ctx["config"])
DB:
-- 每 5 秒触发一次
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("mini_job", "interval", '{"seconds": 5}', 1, 1)
ON DUPLICATE KEY UPDATE enabled=1, schedule_expr='{"seconds":5}', version=version+1;
示例 2:并发抓取 HTTP(有界并发 + 超时)
场景:你要扫一堆 URL,并发抓取,但要避免瞬间打爆下游(有界并发)。
import asyncio
import httpx
from piko import PikoApp
app = PikoApp("http_sweeper")
api_app = app.api_app
URLS = [
"https://example.com",
"https://www.python.org",
# ...
]
@app.job(job_id="sweep_http")
async def sweep_http(ctx, scheduled_time):
concurrency = 50
sem = asyncio.Semaphore(concurrency)
async with httpx.AsyncClient(timeout=10) as client:
async def fetch(url: str):
async with sem:
r = await client.get(url)
return url, r.status_code, len(r.content)
# 关键点:gather + semaphore = 有界并发
results = await asyncio.gather(*(fetch(u) for u in URLS), return_exceptions=True)
ok = [x for x in results if not isinstance(x, Exception)]
print("ok=", len(ok), "total=", len(results))
要点:
- Semaphore 控制并发度(避免把带宽/连接池/下游打爆)
- httpx 是异步 I/O,
gather会把等待 I/O 的时间“让出”给别的协程
示例 3:并发扇出/扇入(fan-out / fan-in)
场景:一条任务输入 → 并发处理 N 份子任务 → 汇总结果。
import asyncio
from piko import PikoApp
app = PikoApp("fanout")
api_app = app.api_app
@app.job(job_id="fanout_fanin")
async def fanout_fanin(ctx, scheduled_time):
items = list(range(1, 501))
async def work(x: int) -> int:
# 模拟 I/O
await asyncio.sleep(0.01)
return x * x
# 扇出:并发执行
results = await asyncio.gather(*(work(x) for x in items))
# 扇入:汇总
total = sum(results)
print("sum=", total)
示例 4:生产者-消费者(asyncio.Queue 背压)
场景:抓取(快)+ 处理(慢),需要 队列缓冲 + 背压(Queue 有界)。
import asyncio
from piko import PikoApp
app = PikoApp("queue_pipeline")
api_app = app.api_app
@app.job(job_id="producer_consumer")
async def producer_consumer(ctx, scheduled_time):
# 有界队列 = 背压点
q: asyncio.Queue[int] = asyncio.Queue(maxsize=200)
concurrency = 20
async def producer():
for i in range(5000):
# 队列满会阻塞(异步阻塞,不占线程)
await q.put(i)
for _ in range(concurrency):
# 结束信号
await q.put(-1)
async def consumer(worker_id: int):
processed = 0
while True:
x = await q.get()
try:
if x == -1:
return processed
# 模拟 I/O 或业务处理
await asyncio.sleep(0.002)
processed += 1
finally:
q.task_done()
consumers = [asyncio.create_task(consumer(i)) for i in range(concurrency)]
prod_task = asyncio.create_task(producer())
await prod_task
# 等待队列处理完
await q.join()
stats = await asyncio.gather(*consumers)
print("total_processed=", sum(stats))
要点:
Queue(maxsize=N)= 背压:生产太快会自动阻塞,防止 OOMq.join()+task_done()= 可靠等待“处理完”
示例 5:异步 I/O + CPU 混合流水线(ProcessPool MapReduce)
场景:先异步下载/读取数据,再做 CPU 重计算(比如解压、解析、特征提取)。
Piko 内置 CpuManager(多进程):
import math
from piko import PikoApp
app = PikoApp("io_cpu_mix")
api_app = app.api_app
def heavy_cpu(x: int) -> int:
# CPU 密集:会占满 GIL(所以要多进程)
math.factorial(2000)
return x * x
@app.job(job_id="io_plus_cpu")
async def io_plus_cpu(ctx, scheduled_time):
items = list(range(1000))
# MapReduce:在多个子进程并行执行 heavy_cpu
results = await app.cpu_manager.map_reduce(
map_fn=heavy_cpu,
items=items,
# 控制并行进程任务数
concurrency=4,
)
print("done:", len(results), "sample:", results[0])
这个场景经常遇到:异步 I/O 把数据拉回来 → 多进程做 CPU 重活 → 异步写出去。
示例 6:把同步阻塞库变成“异步可并发”
场景:你依赖一个同步 SDK(例如某些老库/驱动/算法),但你想在 asyncio 下并发调用。
两种办法:
6.1 用 asyncio.to_thread(适合 I/O 或轻 CPU)
import asyncio
from piko import PikoApp
app = PikoApp("sync_to_async")
api_app = app.api_app
def blocking_call(x: int) -> int:
# 模拟同步阻塞
import time
time.sleep(0.05)
return x + 1
@app.job(job_id="to_thread_demo")
async def to_thread_demo(ctx, scheduled_time):
sem = asyncio.Semaphore(100)
async def run_one(x: int):
async with sem:
return await asyncio.to_thread(blocking_call, x)
results = await asyncio.gather(*(run_one(i) for i in range(1000)))
print("done", len(results))
6.2 用 app.cpu_manager.submit(适合 CPU 重活,需要绕过 GIL)
from piko import PikoApp
app = PikoApp("cpu_submit")
api_app = app.api_app
def cpu_heavy(x: int) -> int:
import math
math.factorial(3000)
return x * 2
@app.job(job_id="cpu_submit_demo")
async def cpu_submit_demo(ctx, scheduled_time):
res = await app.cpu_manager.submit(cpu_heavy, 21)
print("res=", res)
示例 7:有状态任务(水位线 + 自动补跑)
场景:每小时同步一次数据,服务停机 3 小时后重启,要补齐漏掉的数据窗口。
Piko 通过以下手段来实现:
@app.job(..., stateful=True, backfill_policy=...)scheduled_job.last_data_time水位线(成功后自动更新)
from piko import PikoApp
from piko.core.types import BackfillPolicy, DataInterval
app = PikoApp("stateful")
api_app = app.api_app
@app.job(job_id="sync_orders", stateful=True, backfill_policy=BackfillPolicy.CATCH_UP)
async def sync_orders(ctx, scheduled_time):
interval: DataInterval = ctx["data_interval"]
print("sync window:", interval.start, "->", interval.end)
# 你的增量逻辑:WHERE updated_at >= start AND updated_at < end
# await do_incremental_sync(interval.start, interval.end)
调度(每小时一次):
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("sync_orders", "cron", '{"minute": 0}', 1, 1)
ON DUPLICATE KEY UPDATE enabled=1, schedule_expr='{"minute":0}', version=version+1;
补跑策略说明:
BackfillPolicy.CATCH_UP:补齐所有漏掉的窗口(数据完整性优先)BackfillPolicy.SKIP:只跑最新窗口(实时性优先)
示例 8:资源依赖注入 Resource(连接池/客户端自动释放)
Resource 的本质:一个异步上下文管理器工厂。 Piko 在每次 job 执行时会用 AsyncExitStack 自动 enter/exit,确保资源释放。
8.1 定义一个 Resource(例如 HTTP Client)
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
import httpx
from piko.core.resource import resource
@resource(name="http_client")
@asynccontextmanager
async def http_client(ctx: dict[str, object]) -> AsyncGenerator[httpx.AsyncClient, None]:
async with httpx.AsyncClient(timeout=10) as client:
yield client
8.2 在 job 里声明并注入
import asyncio
from piko import PikoApp
app = PikoApp("resource_di")
api_app = app.api_app
@app.job(
job_id="fetch_with_resource",
resources={"client": http_client},
)
async def fetch_with_resource(ctx, scheduled_time, client):
# client 是被注入的 HTTP 客户端
urls = ["https://example.com"] * 100
sem = asyncio.Semaphore(50)
async def fetch(url):
async with sem:
r = await client.get(url)
return r.status_code
codes = await asyncio.gather(*(fetch(u) for u in urls))
print("200_count=", sum(1 for c in codes if c == 200))
示例 9:自定义 Resource(你想注入什么都行)
下面是三个最常见的 Resource 形态:连接池 / SDK Client / 共享缓存。
9.1 注入 Redis(示意)
from contextlib import asynccontextmanager
from piko.core.resource import resource
@resource(name="redis")
@asynccontextmanager
async def redis(ctx: dict[str, object]):
# 在这里创建客户端,并在 finally 中释放连接
client = create_redis_client("redis://127.0.0.1:6379/0")
try:
yield client
finally:
await client.close()
9.2 注入 Mongo / ES / MQ / 任何你自己的 Client
你只要保证:
acquire(ctx)返回一个 async contextmanageryield出你希望注入到 job 的实例finally里把连接关闭/释放即可
示例 10:持久化写入(队列缓冲 + 批量写 + 磁盘兜底)
Piko 的 PersistenceWriter 是一个“生产者-消费者”写入引擎:
- job 里 enqueue(快)
- writer 后台批量 flush 到 sink(可控)
- 写失败会 dump 到磁盘,启动时自动恢复
10.1 写一个最小 Sink
from piko.persistence.sink_base import ResultSink
from piko.persistence.intent import WriteIntent
class PrintSink(ResultSink):
def __init__(self):
super().__init__(name="print")
async def write_batch(self, batch: list[WriteIntent]):
for intent in batch:
print("[SINK]", intent.key, intent.payload)
10.2 注册 Sink,并在 job 中 enqueue
from piko import PikoApp
from piko.persistence.intent import WriteIntent
app = PikoApp("persist_demo")
api_app = app.api_app
# 在应用启动时注册 sink(建议写在 main 模块里)
app.writer.register_sink(PrintSink())
@app.job(job_id="produce_intents")
async def produce_intents(ctx, scheduled_time):
for i in range(1000):
intent = WriteIntent(
sink="print",
key=str(i),
payload={"i": i},
job_id=ctx["job_id"],
run_id=ctx["run_id"],
scheduled_time=scheduled_time,
)
# 队列满会背压阻塞(异步阻塞)
await app.writer.enqueue(intent)
示例 11:TypedSink 类型路由(不同模型不同写法)
如果你的 payload 是不同的 Pydantic Model,希望不同类型走不同写入逻辑:
from pydantic import BaseModel
from piko.persistence.sink_base import TypedSink, on
from piko.persistence.intent import WriteIntent
class User(BaseModel):
id: int
name: str
class Order(BaseModel):
id: int
amount: float
class MyTypedSink(TypedSink):
def __init__(self):
super().__init__(name="typed")
@on(User)
async def write_users(self, users: list[User]):
print("users:", len(users))
@on(Order)
async def write_orders(self, orders: list[Order]):
print("orders:", len(orders))
在 job 里 enqueue:
from piko import PikoApp
from piko.persistence.intent import WriteIntent
app = PikoApp("typed_sink_demo")
api_app = app.api_app
app.writer.register_sink(MyTypedSink())
@app.job(job_id="emit_models")
async def emit_models(ctx, scheduled_time):
await app.writer.enqueue(WriteIntent(
sink="typed",
key="u1",
payload=User(id=1, name="alice"),
job_id=ctx["job_id"],
run_id=ctx["run_id"],
scheduled_time=scheduled_time,
))
await app.writer.enqueue(WriteIntent(
sink="typed",
key="o1",
payload=Order(id=1, amount=9.9),
job_id=ctx["job_id"],
run_id=ctx["run_id"],
scheduled_time=scheduled_time,
))
示例 12:定时任务三种触发器(cron/interval/date)
scheduled_job.schedule_type 支持:
cron:类似 crontabinterval:固定间隔date:单次触发
-- cron:每天 02:30
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("daily_job", "cron", '{"hour": 2, "minute": 30}', 1, 1);
-- interval:每 10 秒
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("fast_job", "interval", '{"seconds": 10}', 1, 1);
-- date:2026-01-08 10:00 触发一次(注意时区由 settings.timezone 决定)
INSERT INTO scheduled_job(job_id, schedule_type, schedule_expr, enabled, version)
VALUES ("one_shot", "date", '{"run_date": "2026-01-08 10:00:00"}', 1, 1);
示例 13:动态配置热更新(灰度生效 effective_from)
在代码里声明 schema:
from pydantic import BaseModel
from piko import PikoApp
app = PikoApp("cfg_demo")
api_app = app.api_app
class SweepConfig(BaseModel):
concurrency: int = 50
timeout_s: float = 10
@app.job(job_id="sweep_cfg", schema=SweepConfig)
async def sweep_cfg(ctx, scheduled_time):
cfg: SweepConfig = ctx["config"] # 已被 Pydantic 校验 & 类型化
print("cfg:", cfg.concurrency, cfg.timeout_s)
在 DB 里更新参数(立即生效):
INSERT INTO job_config(job_id, schema_version, config_json, version)
VALUES ("sweep_cfg", 1, '{"concurrency": 200, "timeout_s": 3}', 1)
ON DUPLICATE KEY UPDATE config_json='{"concurrency":200,"timeout_s":3}', version=version+1;
灰度生效(未来时间生效):
UPDATE job_config
SET config_json='{"concurrency":100,"timeout_s":5}',
effective_from = DATE_ADD(UTC_TIMESTAMP(6), INTERVAL 10 MINUTE),
version=version+1
WHERE job_id="sweep_cfg";
示例 14:多实例部署(Leader Election)
多实例部署时,Piko 默认启用 Leader Election(基于 DB 租约):
- Leader 执行调度与 job run
- Follower 返回
/readyz: standby
配置项(可在 settings.toml 或环境变量中设置):
[default]
leader_enabled = true
leader_name = "default"
leader_lease_s = 30
leader_renew_interval_s = 10
常见部署方式:
- Kubernetes 部署 2~3 个副本
- Prometheus 抓取每个副本的
/metrics - 只有 Leader 的 job_run 会增长(Follower standby)
数据库迁移
Piko 启动时只校验 schema 版本,不会自动执行 DDL——每个应用副本启动时跑 DDL 在多副本部署下是危险的。首次部署或升级时,请在发布流程中单独执行迁移(典型做法是用一个 init-job / 发布前脚本)。
安装即可用的 CLI
从 v0.1.10 起,迁移脚本与 Alembic 配置随 piko 包一起发布,安装后直接通过 piko db 命令管理数据库 schema,无需从源码仓库复制任何文件:
pip install piko-cucc
# 数据库连接通过 PIKO_MYSQL_DSN 配置(与应用启动使用同一配置源)
export PIKO_MYSQL_DSN="mysql+aiomysql://user:pass@host:3306/piko?charset=utf8mb4"
piko db upgrade # 升级到 head(首次部署会创建全部表)
piko db current # 查看当前 schema 版本
piko db downgrade -1 # 回退一步(降级前务必备份数据库)
piko db history # 查看完整迁移版本链
每个子命令都支持 --lock-timeout-s(默认读 PIKO_MIGRATION_LOCK_TIMEOUT_S=30)。
多副本安全:MySQL advisory lock
piko db upgrade/downgrade 在执行 Alembic 期间持有 MySQL 数据库级 advisory lock(锁名 piko-schema-migration),保证同一数据库同时只有一个副本执行 DDL。在 Kubernetes 等多副本环境,多个 init-job 同时启动迁移时,只有一个会真正执行,其余会等待锁(最多 --lock-timeout-s 秒,超时则失败退出)。
Schema 版本与升级路径
当前 schema head 为 schema_v1——表示当前表结构定型的第一版基线。PikoApp 启动时会校验数据库的 alembic_version 必须等于 schema_v1,不匹配则启动失败(避免代码与 schema 不一致导致运行时错误)。
- 首次部署:空库执行
piko db upgrade即可,会按迁移链建出全部表并把版本写到schema_v1。 - 从 v0.1.9 及更早版本升级:旧库的版本号是
0005_remove_sink_dedupe,执行一次piko db upgrade会自动迁移到schema_v1(仅更新版本号,无表结构变更),无损升级。务必先迁移再启动新版本应用。 - 查看当前版本:
piko db current或直接SELECT version_num FROM alembic_version。
离线 SQL(DBA 审核场景)
发布前可生成不连接数据库的迁移 SQL,交给 DBA 审核后手动执行:
# 在仓库内(需要迁移文件)生成离线 SQL
uv run alembic -c piko/migrations/alembic.ini upgrade head --sql > /tmp/piko-upgrade.sql
业务项目自定义表的推荐做法
Piko 只管理自己的 5 张表(scheduled_job/job_config/job_run/job_lock/scheduler_leader)。如果你的业务项目还需要额外的表,不要修改 Piko 的迁移文件(会导致升级时冲突),推荐两种方式:
- 独立 Alembic env(推荐):业务项目使用自己的 Alembic 配置和迁移目录,与 Piko 的迁移链完全解耦。两套迁移可以共用同一个数据库,各自的
alembic_version表互不干扰(业务项目用独立的version_table)。 - 合并 metadata:如果希望统一管理,业务项目的
target_metadata可以同时包含piko.infra.db.Base.metadata和业务自身的 metadata,但 autogenerate 时应排除 Piko 的表(在env.py的include_object钩子里过滤),避免业务迁移反向修改 Piko 表。
create_all_tables() 仅用于测试
piko.infra.db.create_all_tables() 通过 Base.metadata.create_all 直接建表,不写 alembic_version 记录,仅用于隔离测试或本地临时库的快速初始化。生产环境必须通过 piko db upgrade 管理 schema,不要使用此函数。
配置
Piko 使用 Dynaconf,默认读取:
- 包内只读 defaults
PIKO_CONFIG_DIR=/etc/piko指向目录中的piko.tomlPIKO_SETTINGS_PATH=/etc/piko/piko.toml指向单个显式文件PIKO_*环境变量(优先级最高)
默认不读取当前工作目录的 TOML 或 .env。迁移旧部署时可显式设置 PIKO_ENABLE_CWD_CONFIG=true,但该兼容层应尽快移除。显式文件或目录不存在、不可读时启动直接失败;日志只记录来源和覆盖键名,不记录 DSN、密码或 token。
最低必配:
mysql_dsn(必须,负责存元数据)
环境变量示例:
export PIKO_MYSQL_DSN="mysql+aiomysql://user:pass@host:3306/piko?charset=utf8mb4"
export PIKO_TIMEZONE="Asia/Shanghai"
export PIKO_DEBUG="false"
配置契约
除 piko_system_settings.poll_interval_s 外,以下配置均在组件创建或启动时读取,修改后需要重启;PIKO_* 是对应环境变量名。mysql_dsn 是唯一必填项。
| 配置 | 类型 / 默认值 | 环境变量 | 生效域与时机 | 敏感 |
|---|---|---|---|---|
mysql_dsn |
URL / 必填 | PIKO_MYSQL_DSN |
启动静态,建连接池 | 是 |
mysql_pool_size, mysql_max_overflow, mysql_pool_recycle_s |
int / 20, 10, 3600 |
PIKO_MYSQL_* |
启动静态 | 否 |
timezone |
string / Asia/Shanghai |
PIKO_TIMEZONE |
启动静态,Scheduler | 否 |
poll_interval_s, poll_jitter_s |
number / 10, 2 |
PIKO_POLL_* |
poll_interval_s 可由 DB 动态覆盖;poll_jitter_s 重启生效 |
否 |
ap_misfire_grace_s_default, ap_max_instances_default |
int / 300, 1 |
PIKO_AP_* |
启动静态,Scheduler 默认值 | 否 |
leader_enabled, leader_name, leader_lease_s, leader_renew_interval_s |
bool/string/int / true, default, 30, 10 |
PIKO_LEADER_* |
启动静态,Leader | 否 |
leader_watchdog_jitter_s, leader_db_failure_grace_s |
number / 1, 5 |
PIKO_LEADER_* |
Watchdog 运行参数 | 否 |
cpu_workers, per_job_cpu_max |
int / 0, 8 |
PIKO_CPU_*, PIKO_PER_JOB_CPU_MAX |
启动静态,CpuManager | 否 |
job_lock_lease_s, job_lock_heartbeat_s, job_run_orphan_timeout_s |
int / 300, 30, 900 |
PIKO_JOB_* |
Runner 运行参数 | 否 |
job_handler_timeout_s, shutdown_timeout_s |
int / 300, 30 |
PIKO_JOB_HANDLER_TIMEOUT_S, PIKO_SHUTDOWN_TIMEOUT_S |
Runner / 应用停机预算 | 否 |
persist_queue_max, persist_batch_size, persist_batch_timeout_s, persist_flush_timeout_s, persist_disk_fallback_path |
int/number/path / 200, 100, 0.5, 60, /tmp/piko_fallback.bin |
PIKO_PERSIST_* |
Writer 创建或停机时读取 | 路径视部署而定 |
metrics_enabled, api_docs_enabled |
bool / true, false |
PIKO_METRICS_ENABLED, PIKO_API_DOCS_ENABLED |
API 路由创建或请求时 | 否 |
debug, log_level, log_json |
bool/string/bool / false, INFO, false |
PIKO_DEBUG, PIKO_LOG_* |
启动静态,日志初始化 | 否 |
系统级 DB 动态配置只接受下面这一项,其他字段写入 piko_system_settings 不会生效:
{"poll_interval_s": 15}
有效范围是 1 到 3600 秒;非法值会保留上一次有效间隔并记录告警。任务级参数仍写入对应 job_config,调度字段仍写入 scheduled_job。
运维端点与可观测性
Piko 内置 FastAPI 运维端点(无需你自己写):
GET /healthz:存活探针(liveness)GET /readyz:就绪探针(readiness;Follower 会是 standby)GET /metrics:Prometheus 指标
这些端点只应绑定在私网/监控网段,或放在已有认证反向代理和 NetworkPolicy 后面;Piko 不提供默认的 HTTP 认证。生产默认关闭 /docs、/redoc 和 /openapi.json,仅在 api_docs_enabled = true 且访问边界受控时临时开启。/readyz 会在数据库、Writer、Watcher、Scheduler 或 Leader 未就绪时返回 HTTP 503。
常见指标(示意):
- job 成功/失败计数:
JOB_RUN_TOTAL{job_id=..., status=...} - job 耗时直方图:
JOB_DURATION_SECONDS{job_id=...} - leader 状态:
LEADER_STATUS{host=...} - 持久化队列长度:
PERSISTENCE_QUEUE_SIZE
项目结构建议
一个推荐的业务项目结构:
my_project/
app.py # 创建 PikoApp、注册 sink、启动
my_project/
__init__.py
jobs.py # 简单场景:集中放 job
jobs/
__init__.py
user_sync/
__init__.py
jobs.py # 复杂场景:分目录,配合 autodiscover
report/
__init__.py
jobs.py
在 app.py 中:
from piko import PikoApp, autodiscover
app = PikoApp("my_project")
# 自动导入所有 *.jobs.py,触发注册
autodiscover("my_project", module_name="jobs")
api_app = app.api_app
FAQ
1) 我的任务里怎么拿到配置?
- 给 job 声明
schema=YourConfigModel - 然后
cfg: YourConfigModel = ctx["config"]
2) 如何控制并发?
- I/O 并发:
asyncio.Semaphore+gather - 生产者/消费者:
asyncio.Queue(maxsize=...) - CPU 并行:
app.cpu_manager.map_reduce(..., concurrency=N)
3) 如何注入数据库/Redis/Mongo 等资源?
用 Resource:
- 用
@resource(name=...)包装一个@asynccontextmanager函数 - 在资源函数中创建连接/客户端并
yield - 在 job 的
resources={...}中声明资源类
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file piko_cucc-0.1.10.tar.gz.
File metadata
- Download URL: piko_cucc-0.1.10.tar.gz
- Upload date:
- Size: 100.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via:
twine/6.2.0 CPython/3.12.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
2d8c667852c85bde336132960fc86d478e463841209a2a5e39760fd9d89cb9e7
|
|
| MD5 |
2468af65260a5e4d3920c01117e7cc01
|
|
| BLAKE2b-256 |
feb7e10fda07685ff12c1da53b22766da57738ffb1fd84b9b0d63fc2d1bcf021
|
File details
Details for the file piko_cucc-0.1.10-py3-none-any.whl.
File metadata
- Download URL: piko_cucc-0.1.10-py3-none-any.whl
- Upload date:
- Size: 105.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via:
twine/6.2.0 CPython/3.12.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
28180fd485738b2abc1ac9d14b7d58dca31d1c47074441640babb29919416bd4
|
|
| MD5 |
ce1a9dde4ead2630aff5e4f51b66edc7
|
|
| BLAKE2b-256 |
28947cd3716e99e7e27d047d0f5f161c429f5143a9e32fefa1057feea836b929
|