You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Camunda 8 Python 如何定向操作指定流程实例完成任务

实现方法

核心原理

pyzeebe 默认开启任务自动完成机制:只要任务处理函数正常返回,Worker 会自动向 Zeebe 发送任务完成请求,所有被拉取激活的任务都会被标记完成,这就是所有实例任务被自动执行的根本原因。

Job 对象本身已经内置了流程实例唯一标识字段 job.process_instance_key,和启动流程时拿到的 process_instance_key 完全对应,可以直接用来做单实例精准筛选。

具体修改步骤

  1. 关闭对应任务的自动完成开关
    给 task 装饰器加 auto_complete=False 参数,关闭自动完成后,只有主动调用 job.complete() 才会标记任务完成,函数返回值不会再触发自动完成逻辑,完全由你控制任务的完成时机。

  2. 增加实例筛选逻辑
    任务函数拿到 Job 对象后,首先比对 job.process_instance_key 和你需要处理的目标流程实例键:

  • 匹配目标实例:执行对应业务逻辑,手动调用 await job.complete(返回变量字典) 完成该实例的任务
  • 不匹配目标实例:显式调用失败接口,保持任务原有重试次数,设置短延后把任务放回待处理队列,不会标记任务完成,后续轮询还能再次拉取到该任务。
  1. (可选优化)限制Worker单次拉取任务数
    把 Worker 的 max_jobs_to_activate 设为较小值(比如1),避免一次性拉取大量非目标任务占用激活锁,降低Broker无效请求压力。

修改后的可运行代码

import asyncio
import nest_asyncio
from pyzeebe import ZeebeTaskRouter, ZeebeWorker, create_insecure_channel, Job
nest_asyncio.apply()

# 全局配置:指定当前要处理的流程实例key,可根据业务需求动态修改
TARGET_PROCESS_INSTANCE_KEY = 123456789  # 替换为实际需要处理的实例key

async def main():
    channel = create_insecure_channel(hostname="localhost", port=26500)
    # 配置worker每次最多拉1个任务,避免占用过多待处理任务
    worker = ZeebeWorker(channel, max_jobs_to_activate=1)
    router = ZeebeTaskRouter()
    
    async def exception_handler(exception: Exception, job: Job) -> None:
        await job.set_error_status(message='PropertyError', error_code='PropertyError') 

    # 关闭自动完成,挂载异常处理器
    @router.task(task_type="TimerValue", auto_complete=False, exception_handler=exception_handler)
    async def TimerValue(job: Job):
        # 先判断是否为目标流程实例
        if job.process_instance_key == TARGET_PROCESS_INSTANCE_KEY:
            # 是目标实例,执行业务逻辑后手动完成任务
            print(f"处理目标实例任务,process_instance_key: {job.process_instance_key}")
            await job.complete({"timer_value": "P14D"})
        else:
            # 非目标实例,保持重试次数不变,延后1秒放回待处理队列
            print(f"跳过非目标实例任务,process_instance_key: {job.process_instance_key}")
            await job.fail(remaining_retries=job.retries, retry_backoff_ms=1000)
    
    worker.include_router(router)
    await worker.work()

if __name__ == "__main__":
    asyncio.run(main())
    print("done")

注意事项

  • 非目标任务调用fail接口时一定要保持remaining_retries和当前job的重试次数一致,否则重试次数耗尽后任务会触发 incident 无法继续处理
  • retry_backoff_ms可以根据业务场景调整,不需要秒级响应的场景可以设为5000或更高,减少无效轮询的性能消耗
  • 如果需要动态切换要处理的目标实例,直接修改全局TARGET_PROCESS_INSTANCE_KEY的值即可,不需要重启Worker

内容的提问来源于stack exchange,提问作者Sachin Rajput

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 23:24:14