Skip to content

feat(periodic): 新增主备定时任务 SDK 与调用示例 - #41

Open
openjiuwen-sync-bot[bot] wants to merge 1 commit into
openJiuwen-ai:developfrom
openjiuwenai:sync/pr-417
Open

feat(periodic): 新增主备定时任务 SDK 与调用示例#41
openjiuwen-sync-bot[bot] wants to merge 1 commit into
openJiuwen-ai:developfrom
openjiuwenai:sync/pr-417

Conversation

@openjiuwen-sync-bot

@openjiuwen-sync-bot openjiuwen-sync-bot Bot commented Aug 13, 2026

Copy link
Copy Markdown

Paired: GitHub #41GitCode !417

What type of PR is this?

/kind feature
Self-checklist:(请自检,在[ ]内打上x,我们将检视你的完成情况,否则会导致pr无法合入

    • 设计:PR对应的方案是否已经经过Maintainer评审,方案检视意见是否均已答复并完成方案修改
    • 测试:PR中的代码是否已有UT/ST测试用例进行充分的覆盖,新增测试用例是否随本PR一并上库或已经上库
    • 验证:PR描述信息中是否已包含对该PR对应的Feature、Refactor、Bugfix的预期目标达成情况的详细验证结果描述
    • 接口:是否涉及对外接口变更,相应变更已得到接口评审组织的通过,API对应的注释信息已经刷新正确
    • 文档:是否涉及官网文档修改,如果涉及请及时提交资料到Doc仓

Summary

新增进程内主备(SingleLeader)周期任务 SDKfoundation.periodic),对外通过工厂 create_single_leader_job 创建任务;并提供调用示例 applications/periodic_job/example.py
核心能力:

  • 固定间隔调度(IntervalSchedule
  • 开火前集合窗口报名 + Redis Lua 抽签选主
  • 执行锁 SET NX EX + 持锁续期 + token 校验释放
  • 同一时刻全局仅一个实例执行 on_tick(无参回调)
  • 已在多台服务器上测试,多实例获取锁的概率相同。
  • 示例代码放在applications\periodic_job\example.py。

Co-authored-by: zhangxiangyu52 <zhangxiangyu52@huawei.com>
@CLAassistant

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.
You have signed the CLA already but the status is still pending? Let us recheck it.

@openjiuwen-collaboration-bot

Copy link
Copy Markdown

head_sha: b772e455905ce072f7bc00174719a0d951bf07ab

变更摘要

本 PR 新增进程内主备(SingleLeader)周期任务 SDKopenjiuwen_runtime.foundation.periodic),通过工厂函数 create_single_leader_job 组装调度、协调器与执行器,并提供调用示例 applications/periodic_job/example.py。核心能力包括固定间隔调度(IntervalSchedule)、开火前集合窗口报名并借助 Redis Lua 抽签选主(SingleLeaderCoordinator)、执行锁 SET NX EX 加持锁续期与 token 校验释放(TickLock),以及由 JobRunner 统一管理任务生命周期与到点回调,保证同一时刻全局仅一个实例执行无参的 on_tick

主要改动

  • 新增 SDK 对外入口与工厂函数:在 factory.py 中定义 create_single_leader_job,以 nameon_tickinstance_id 为必填参数组装 JobRunnerIntervalScheduleSingleLeaderCoordinator,并在 periodic/__init__.py 导出对外 API。
  • 实现主备抽签协调器SingleLeaderCoordinator 在开火窗口内通过 SADD 报名、到期后用 Lua SRANDMEMBER + SET NX winner 原子选主,仅胜者继续获取执行锁并启动续期。
  • 实现执行锁与续期机制TickLock 使用 SET NX EX 抢锁,通过 Lua 实现仅 owner 可续期(EXPIRE)与仅 owner 可释放(DEL),并支持后台 asyncio 续期任务和 token 校验。
  • 新增任务运行器与间隔调度JobRunner 负责提前窗口醒来、协调选主、执行回调并放锁,支持 start/stoprun_on_start 与重叠跳过;IntervalSchedule 按整秒/整 N 秒边界对齐计算下次触发时间。
  • 新增调用示例applications/periodic_job/example.py 演示了如何注入 Redis、创建主备任务并 start/stop,并说明各选填参数(如 interval_secgather_window_seclock_keyrun_on_start)的用途。

@openjiuwen-collaboration-bot

openjiuwen-collaboration-bot Bot commented Aug 13, 2026

Copy link
Copy Markdown

head_sha: b772e455905ce072f7bc00174719a0d951bf07ab

代码审查

审查总结

本次审查覆盖了全部 11 个变更文件,结论如下:

发现的问题按优先级统计:

  • P2:5 个(锁续期任务异常静默退出、失锁状态无人消费导致脑裂、候选集合提前过期导致任务静默失效、stop/start 竞态产生双循环、run_on_start 首轮 tick 异常终结任务)
  • P3:2 个(expire 失败被 debug 吞掉导致 Redis 键泄漏、_busy 重叠跳过为死代码)
  • P0/P1:0 个

总体风险判断: 该 PR 的核心主备选主 + 锁续期机制方向正确、Lua 脚本参数化无注入风险、导入路径与现有 openjiuwen_runtime.foundation.* 约定一致。但锁生命周期与失锁处理存在多处可靠性缺口——尤其"失锁后不停止执行"与"续期任务异常静默退出"会直接削弱"同一时刻全局仅一个实例执行"这一核心保证;此外 run_on_start 路径的异常处理不对称会导致任务在启动时被意外终结。建议在合入前优先修复这些 P2 问题。

逐文件审查确认:

  • applications/periodic_job/example.py — 无问题(示例代码,导入与用法正确)
  • foundation/openjiuwen_runtime/foundation/periodic/__init__.py — 无问题(导出齐全)
  • foundation/openjiuwen_runtime/foundation/periodic/coordinator/__init__.py — 无问题
  • foundation/openjiuwen_runtime/foundation/periodic/coordinator/base.py — 无问题(Protocol 定义)
  • foundation/openjiuwen_runtime/foundation/periodic/coordinator/single_leader.py — 发现 2 个问题(候选集合提前过期、expire 失败吞异常)
  • foundation/openjiuwen_runtime/foundation/periodic/factory.py — 无独立问题(gather_window_sec 无上限约束已并入 single_leader 相关发现)
  • foundation/openjiuwen_runtime/foundation/periodic/lock.py — 发现 2 个问题(续期异常静默退出、失锁状态无人消费)
  • foundation/openjiuwen_runtime/foundation/periodic/runner.py — 发现 3 个问题(stop/start 竞态、_busy 死代码、run_on_start 异常处理缺失)
  • foundation/openjiuwen_runtime/foundation/periodic/schedule/__init__.py — 无问题
  • foundation/openjiuwen_runtime/foundation/periodic/schedule/base.py — 无问题(Protocol 定义)
  • foundation/openjiuwen_runtime/foundation/periodic/schedule/interval.py — 无问题(边界对齐逻辑正确,始终返回严格大于 now 的值)
类型 数量
🔴 阻塞 0
🟡 建议 2

💬 仅评论

@openjiuwen-collaboration-bot

Copy link
Copy Markdown

head_sha: b772e455905ce072f7bc00174719a0d951bf07ab

任务名称 结果 日志操作
静态检查 ❌FAILED 点此跳转
禁用词扫描 N/A N/A
防投毒检查 ✅SUCCESS 点此跳转
开源合规检查 ✅SUCCESS 点此跳转
UT测试 N/A N/A
ST测试 N/A N/A
build 编译包 N/A N/A
ruff codecheck {{ruff_codecheck_result}} {{ruff_codecheck_link_text}}


async def _run_forever(self) -> None:
if self._run_on_start and not self._stopped.is_set():
await self._safe_tick(planned_fire=self._now())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

head_sha: b772e455905ce072f7bc00174719a0d951bf07ab

🟡 Medium Priority

建议:将 run_on_start 的首轮 tick 包进 try/except Exception,记录日志后继续进入主循环,保证与主循环一致的容错。

Suggested change
await self._safe_tick(planned_fire=self._now())
if self._run_on_start and not self._stopped.is_set():
try:
await self._safe_tick(planned_fire=self._now())
except Exception:
logger.exception(
"JobRunner run_on_start failed: job=%s instance=%s",
self._name,
self._instance_id,
)

return
logger.debug("tick lock renewed: key=%s", self._lock_key)
except asyncio.CancelledError:
return

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

head_sha: b772e455905ce072f7bc00174719a0d951bf07ab

🟡 Medium Priority

建议:在 _renew_loop 内对 renew_once 的异常做处理,失锁或异常都置 _lost=True 并退出,避免续期任务静默死亡。

Suggested change
return
async def _renew_loop(self, token: str) -> None:
try:
while True:
await asyncio.sleep(self._renew_interval_sec)
try:
ok = await self.renew_once(token)
except Exception:
self._lost = True
logger.exception(
"tick lock renew error: key=%s token=%s",
self._lock_key,
token,
)
return
if not ok:
self._lost = True
logger.warning(
"tick lock lost on renew: key=%s token=%s",
self._lock_key,
token,
)
return
logger.debug("tick lock renewed: key=%s", self._lock_key)
except asyncio.CancelledError:
return

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants