Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 52 additions & 0 deletions applications/periodic_job/example.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

"""主备定时任务 —— 怎么调用。

redis 由调用方注入(生产用真 Redis;测试可注入 FakeRedis)。
"""

from __future__ import annotations

import asyncio
import socket
from typing import Any

from openjiuwen_runtime.foundation.periodic import create_single_leader_job


def _instance_id() -> str:
"""本机工号:集群里每台机器必须不一样,否则抽签分不清谁是谁。"""
return socket.gethostname()


async def example_create_single_leader_job(redis: Any) -> None:
async def on_tick() -> None:
# 到点干活;需要时间就自己 time.time()
print("[demo] tick")

job = create_single_leader_job(
redis, # 【必填】Redis(报名 / 抽签 / 执行锁)
name="demo", # 【必填】任务名;不传 lock_key 时 → 锁名 lock:demo
on_tick=on_tick, # 【必填】到点回调(无参)
instance_id=_instance_id(), # 【必填】本机实例 ID
# interval_sec=1, # 【选填】默认 1;执行锁 TTL = 此值
# gather_window_sec=0.08, # 【选填】默认 0.08;开火前提前醒来报名
# lock_key="", # 【选填】空则 lock:{name};一般不用改
# run_on_start=False, # 【选填】True=启动立刻跑一轮(多用于测试)
)
await job.start()
try:
await asyncio.sleep(5) # 演示跑几秒;生产里挂在服务生命周期上
finally:
await job.stop()


if __name__ == "__main__":
# 注入 redis 后跑:
# import redis.asyncio as redis
# r = redis.from_url("redis://127.0.0.1:6379/0")
# asyncio.run(example_create_single_leader_job(r))
raise SystemExit(
"请注入 redis 后调用:asyncio.run(example_create_single_leader_job(your_redis))"
)
26 changes: 26 additions & 0 deletions foundation/openjiuwen_runtime/foundation/periodic/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

"""进程内周期任务 SDK(Schedule + Coordinator + JobRunner)。

主备模式:``SingleLeaderCoordinator``(等待窗口抽签 + 锁续期)。
时间:本机墙钟 ``time.time``。

对外入口:``create_single_leader_job``(唯一配置面)。
"""

from .coordinator import Coordinator, SingleLeaderCoordinator
from .factory import create_single_leader_job
from .lock import TickLock
from .runner import JobRunner
from .schedule import IntervalSchedule, Schedule

__all__ = (
"Coordinator",
"IntervalSchedule",
"JobRunner",
"Schedule",
"SingleLeaderCoordinator",
"TickLock",
"create_single_leader_job",
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

from .base import Coordinator
from .single_leader import SingleLeaderCoordinator

__all__ = (
"Coordinator",
"SingleLeaderCoordinator",
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

"""协调器协议。"""

from __future__ import annotations

from typing import Optional, Protocol


class Coordinator(Protocol):
async def try_claim(
self,
*,
now: float,
instance_id: str,
planned_fire: float | None = None,
) -> Optional[str]:
"""试着领取本轮执行权;成功返回锁 token,失败返回 None。

``planned_fire``:本拍语义上的开火整点(如 10.000)。
提前醒来时 ``now`` 可能是 T-窗口,epoch / 等到点应以 ``planned_fire`` 为准。
"""
...

async def release(self, token: str) -> None:
"""交回执行权(按 token 校验后放锁)。"""
...
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

"""主备协调:提前报名 → 等到开火点 → 抽签选主 → 持锁续期执行。

流程(每拍,配合 JobRunner 提前 ``gather_window`` 醒来):
1. ``planned_fire`` 为本拍整点 T;``now`` 多为 T-窗口
2. ``SADD candidates:{epoch}`` 报名(epoch 取自 T)
3. 睡到 T(剩余窗口),让网络慢的实例也能进来
4. Lua 原子抽签:``SRANDMEMBER`` + ``SET NX winner:{epoch}``
5. 只有 winner 去 ``SET NX`` 执行锁,并启动续期;别人空转
"""

from __future__ import annotations

import asyncio
from typing import Any, Optional

from openjiuwen_runtime.foundation.log import get_logger
from openjiuwen_runtime.foundation.periodic.lock import TickLock

logger = get_logger(__name__)

_ELECT_LUA = """
local existing = redis.call('GET', KEYS[1])
if existing then
return existing
end
local pick = redis.call('SRANDMEMBER', KEYS[2])
if not pick then
return false
end
local ok = redis.call('SET', KEYS[1], pick, 'NX', 'EX', tonumber(ARGV[1]))
if ok then
return pick
end
return redis.call('GET', KEYS[1])
"""


class SingleLeaderCoordinator:
"""主备:开火前窗口内集齐候选人,到点后随机选唯一执行者。"""

def __init__(
self,
redis: Any,
*,
lock_key: str,
lock_ttl_sec: int = 1,
token_prefix: str = "job",
instance_id: str = "",
gather_window_sec: float = 0.08,
meta_ttl_sec: int = 3,
) -> None:
self._redis = redis
self._instance_id = instance_id
self._lock_key = lock_key
self._gather_window_sec = max(float(gather_window_sec), 0.0)
self._meta_ttl_sec = max(int(meta_ttl_sec), 1)
self._lock = TickLock(
redis,
lock_key=lock_key,
lock_ttl_sec=lock_ttl_sec,
token_prefix=token_prefix,
instance_id=instance_id,
)

def _candidates_key(self, epoch: int) -> str:
return f"{self._lock_key}:candidates:{epoch}"

def _winner_key(self, epoch: int) -> str:
return f"{self._lock_key}:winner:{epoch}"

async def try_claim(
self,
*,
now: float,
instance_id: str,
planned_fire: float | None = None,
) -> Optional[str]:
iid = instance_id or self._instance_id
fire_at = float(planned_fire) if planned_fire is not None else float(now)
epoch = int(fire_at)
cand_key = self._candidates_key(epoch)
winner_key = self._winner_key(epoch)

await self._redis.sadd(cand_key, iid)
try:
await self._redis.expire(cand_key, self._meta_ttl_sec)
except Exception:
logger.debug("candidates expire failed: key=%s", cand_key)

if planned_fire is not None:
delay = fire_at - now
else:
delay = self._gather_window_sec
if delay > 0:
await asyncio.sleep(delay)

winner = await self._elect(winner_key, cand_key)
if winner is None:
logger.debug("no candidates for epoch=%s instance=%s", epoch, iid)
return None

winner_s = winner.decode() if isinstance(winner, (bytes, bytearray)) else str(winner)
if winner_s != iid:
logger.debug(
"not elected: epoch=%s instance=%s winner=%s",
epoch,
iid,
winner_s,
)
return None

token = await self._lock.try_acquire()
if token is None:
logger.warning(
"elected but lock busy: epoch=%s instance=%s key=%s",
epoch,
iid,
self._lock_key,
)
return None

self._lock.start_renew(token)
logger.info(
"single_leader claimed: epoch=%s instance=%s key=%s",
epoch,
iid,
self._lock_key,
)
return token

async def _elect(self, winner_key: str, cand_key: str) -> Any:
return await self._redis.eval(
_ELECT_LUA,
2,
winner_key,
cand_key,
str(self._meta_ttl_sec),
)

async def release(self, token: str) -> None:
try:
await self._lock.release_if_owner(token)
except Exception:
logger.exception(
"single_leader release failed: key=%s token=%s",
self._lock.lock_key,
token,
)
60 changes: 60 additions & 0 deletions foundation/openjiuwen_runtime/foundation/periodic/factory.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
# coding: utf-8
# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved

"""工厂:一条调用组装好主备定时任务。

对外配置只认本函数参数;内部零件(Runner / Schedule / Coordinator)不另搞 Config 袋。
"""

from __future__ import annotations

from typing import Any, Awaitable, Callable

from openjiuwen_runtime.foundation.periodic.coordinator.single_leader import (
SingleLeaderCoordinator,
)
from openjiuwen_runtime.foundation.periodic.runner import JobRunner
from openjiuwen_runtime.foundation.periodic.schedule.interval import IntervalSchedule

# 内部常量(不对外暴露)
_META_TTL_SEC = 3


def create_single_leader_job(
redis: Any, # Redis 客户端(要能 async:set/eval/sadd…)
*,
name: str, # 任务名;默认锁 key 为 lock:{name}
on_tick: Callable[[], Awaitable[None]], # 到点回调:async def on_tick() -> None
instance_id: str, # 本机实例 ID,报名/抽签用,集群内需唯一
interval_sec: int = 1, # 每隔多少秒响一次;锁 TTL 与此相同
gather_window_sec: float = 0.08, # 开火前集合窗口:提前醒来报名,到整秒抽签
lock_key: str = "", # 执行锁 Redis key;空则用 lock:{name}
run_on_start: bool = False, # True 启动后立刻跑一轮(一般仅测试)
) -> JobRunner:
"""创建主备周期任务,返回可 start/stop 的 JobRunner。

调用方通常只需:
job = create_single_leader_job(redis, name="x", on_tick=..., instance_id="n1")
await job.start()
...
await job.stop()
"""
interval = max(int(interval_sec), 1)
key = (lock_key or f"lock:{name}").rstrip(":")
return JobRunner(
name=name,
schedule=IntervalSchedule(interval),
coordinator=SingleLeaderCoordinator(
redis,
lock_key=key,
lock_ttl_sec=interval,
token_prefix=name,
instance_id=instance_id,
gather_window_sec=gather_window_sec,
meta_ttl_sec=_META_TTL_SEC,
),
on_tick=on_tick,
instance_id=instance_id,
gather_window_sec=gather_window_sec,
run_on_start=run_on_start,
)
Loading