Skip to main content

Piko - Data-Oriented Async Task Orchestrator

PyPI version Python versions License: MIT GitHub release

一个面向 数据任务(ETL / 同步 / 扫描 / 归档 / 监控) 的异步并发编排框架:

代码里写任务(Job),数据库里配调度与参数(Schedule/Config),运行时自动热更新,天然支持异步高并发与可观测性。


目录


为什么是 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 主要做两件事:

  1. 同步 job 参数(job_config → ConfigCache)

    • job_config.config_json 写入内存缓存(ConfigCache)
    • 支持 effective_from:如果配置的生效时间在未来,则暂不写入缓存(灰度/延时生效)
    • 任务真正执行时会从 ConfigCache 取 ctx["config"],因此 下一次触发执行就会使用新参数
  2. 同步 job 调度(scheduled_job → APScheduler)

    • 对比 DB 中 enabled 的 job 列表与 APScheduler 当前 job 列表:增、删、改
    • 重要:Piko 判断“是否需要更新触发器”的依据是 scheduled_job.version
      • watcher 会把 version 写进 APScheduler job 的 name(例如 v1 / v2
      • 只有当版本标签变化时才会 reschedule_job(...)
      • 这意味着:仅修改 schedule_expr 等配置,而不 bump version,调度不会热更新

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"]:任务 ID
    • ctx["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) = 背压:生产太快会自动阻塞,防止 OOM
  • q.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 contextmanager
  • yield 出你希望注入到 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:类似 crontab
  • interval:固定间隔
  • 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 的迁移文件(会导致升级时冲突),推荐两种方式:

  1. 独立 Alembic env(推荐):业务项目使用自己的 Alembic 配置和迁移目录,与 Piko 的迁移链完全解耦。两套迁移可以共用同一个数据库,各自的 alembic_version 表互不干扰(业务项目用独立的 version_table)。
  2. 合并 metadata:如果希望统一管理,业务项目的 target_metadata 可以同时包含 piko.infra.db.Base.metadata 和业务自身的 metadata,但 autogenerate 时应排除 Piko 的表(在 env.pyinclude_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.toml
  • PIKO_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}

有效范围是 13600 秒;非法值会保留上一次有效间隔并记录告警。任务级参数仍写入对应 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

piko_cucc-0.1.10.tar.gz (100.0 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

piko_cucc-0.1.10-py3-none-any.whl (105.7 kB view details)

Uploaded Python 3

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

Hashes for piko_cucc-0.1.10.tar.gz
Algorithm Hash digest
SHA256 2d8c667852c85bde336132960fc86d478e463841209a2a5e39760fd9d89cb9e7
MD5 2468af65260a5e4d3920c01117e7cc01
BLAKE2b-256 feb7e10fda07685ff12c1da53b22766da57738ffb1fd84b9b0d63fc2d1bcf021

See more details on using hashes here.

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

Hashes for piko_cucc-0.1.10-py3-none-any.whl
Algorithm Hash digest
SHA256 28180fd485738b2abc1ac9d14b7d58dca31d1c47074441640babb29919416bd4
MD5 ce1a9dde4ead2630aff5e4f51b66edc7
BLAKE2b-256 28947cd3716e99e7e27d047d0f5f161c429f5143a9e32fefa1057feea836b929

See more details on using hashes here.

Release history Release notifications | RSS feed

0.1.12

2 files

0.1.11

2 files

This release

0.1.10 This release

2 files

0.1.9

2 files

0.1.8

2 files

0.1.7

2 files

0.1.6

2 files

0.1.5

2 files

0.1.4

2 files

0.1.3

2 files

0.1.2

2 files

0.1.1

2 files

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page