返回博客

审批工作流:Temporal + 状态机的企业级审批引擎

企业审批流程是决策引擎的核心执行环节:一个决策做出后,可能需要经过多级审批才能生效。coomia-dip 使用 Temporal 作为工作流编排引擎,结合 有限状态机(FSM) 管理审批状态,实现了支持串行、并行、条件分支和超时升级的企业级审批引擎。本文深入解析审批工作流的架构设计、Temporal Workflow/Activity 实现、状态机模型、以及与 DecisionEngine 的集成。

Coomia发布于 2025年9月1日13 分钟阅读
分享本文Twitter / X

系列:S5 智能决策 · 第 10 篇 | 难度:高级 | 阅读时间:20 分钟

审批工作流:Temporal + 状态机的企业级审批引擎

#TL;DR

企业审批流程是决策引擎的核心执行环节:一个决策做出后,可能需要经过多级审批才能生效。coomia-dip 使用 Temporal 作为工作流编排引擎,结合 有限状态机(FSM) 管理审批状态,实现了支持串行、并行、条件分支和超时升级的企业级审批引擎。本文深入解析审批工作流的架构设计、Temporal Workflow/Activity 实现、状态机模型、以及与 DecisionEngine 的集成。

#1. 企业审批的复杂性

#1.1 审批场景分类

Code
企业审批场景:

  简单审批                多级审批              条件分支审批
  ┌──────┐              ┌──────┐             ┌──────┐
  │申请人 │──→ 审批人     │申请人 │──→ 主管     │申请人 │──→ 金额判断
  └──────┘              └──┬───┘   ──→ 经理    └──┬───┘
                           │       ──→ VP          │
                           ▼                       ▼
                        ┌──────┐           ┌──────────────┐
                        │ 会签  │           │ <10万: 主管    │
                        │(并行) │           │ 10-50万: 经理  │
                        └──────┘           │ >50万: VP+CFO │
                                           └──────────────┘

#1.2 审批流程的技术挑战

挑战描述解决方案
持久性审批可能持续数天/周Temporal 持久化执行
超时审批人不响应自动提醒 + 超时升级
并行多人会签Temporal 并行 Activity
回退驳回后重新提交状态机循环
可见性实时查看进度Temporal Query
审计完整操作记录Event History

#2. 审批状态机

#2.1 状态模型

Python
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from typing import Any


class ApprovalStatus(Enum):
    """审批状态"""
    DRAFT = "draft"
    PENDING = "pending"
    IN_REVIEW = "in_review"
    APPROVED = "approved"
    REJECTED = "rejected"
    WITHDRAWN = "withdrawn"
    ESCALATED = "escalated"
    EXPIRED = "expired"


class ApprovalAction(Enum):
    """审批动作"""
    SUBMIT = "submit"
    APPROVE = "approve"
    REJECT = "reject"
    WITHDRAW = "withdraw"
    ESCALATE = "escalate"
    RETURN = "return"          # 退回修改
    COUNTERSIGN = "countersign"  # 加签
    TIMEOUT = "timeout"


@dataclass
class ApprovalRequest:
    """审批请求"""
    request_id: str
    decision_id: str             # 关联的决策 ID
    applicant: str
    domain: str
    title: str
    content: dict[str, Any]
    amount: float = 0.0
    priority: str = "normal"     # low, normal, high, urgent
    status: ApprovalStatus = ApprovalStatus.DRAFT
    current_level: int = 0
    approval_chain: list[ApprovalLevel] = field(default_factory=list)
    history: list[ApprovalEvent] = field(default_factory=list)
    created_at: datetime = field(default_factory=datetime.utcnow)
    updated_at: datetime = field(default_factory=datetime.utcnow)


@dataclass
class ApprovalLevel:
    """审批层级"""
    level: int
    approvers: list[str]
    mode: str = "any"            # any=或签, all=会签
    timeout_hours: int = 48
    auto_escalate_to: str | None = None


@dataclass
class ApprovalEvent:
    """审批事件"""
    event_id: str
    action: ApprovalAction
    actor: str
    comment: str = ""
    timestamp: datetime = field(default_factory=datetime.utcnow)
    metadata: dict[str, Any] = field(default_factory=dict)

#2.2 状态转换引擎

Python
class ApprovalStateMachine:
    """审批状态机"""

    # 状态转换表
    _transitions: dict[tuple[ApprovalStatus, ApprovalAction], ApprovalStatus] = {
        (ApprovalStatus.DRAFT, ApprovalAction.SUBMIT): ApprovalStatus.PENDING,
        (ApprovalStatus.PENDING, ApprovalAction.APPROVE): ApprovalStatus.IN_REVIEW,
        (ApprovalStatus.PENDING, ApprovalAction.REJECT): ApprovalStatus.REJECTED,
        (ApprovalStatus.PENDING, ApprovalAction.WITHDRAW): ApprovalStatus.WITHDRAWN,
        (ApprovalStatus.PENDING, ApprovalAction.ESCALATE): ApprovalStatus.ESCALATED,
        (ApprovalStatus.PENDING, ApprovalAction.TIMEOUT): ApprovalStatus.ESCALATED,
        (ApprovalStatus.PENDING, ApprovalAction.RETURN): ApprovalStatus.DRAFT,
        (ApprovalStatus.IN_REVIEW, ApprovalAction.APPROVE): ApprovalStatus.APPROVED,
        (ApprovalStatus.IN_REVIEW, ApprovalAction.REJECT): ApprovalStatus.REJECTED,
        (ApprovalStatus.IN_REVIEW, ApprovalAction.ESCALATE): ApprovalStatus.ESCALATED,
        (ApprovalStatus.ESCALATED, ApprovalAction.APPROVE): ApprovalStatus.APPROVED,
        (ApprovalStatus.ESCALATED, ApprovalAction.REJECT): ApprovalStatus.REJECTED,
        (ApprovalStatus.REJECTED, ApprovalAction.SUBMIT): ApprovalStatus.PENDING,
    }

    def can_transition(self, current: ApprovalStatus,
                       action: ApprovalAction) -> bool:
        return (current, action) in self._transitions

    def transition(self, request: ApprovalRequest,
                   action: ApprovalAction,
                   actor: str,
                   comment: str = "") -> ApprovalStatus:
        key = (request.status, action)
        if key not in self._transitions:
            raise ValueError(
                f"Invalid transition: {request.status.value} + {action.value}"
            )

        new_status = self._transitions[key]

        # 特殊处理:IN_REVIEW 中的 APPROVE 需要检查是否还有更多层级
        if (action == ApprovalAction.APPROVE and
                request.current_level < len(request.approval_chain) - 1):
            request.current_level += 1
            new_status = ApprovalStatus.PENDING

        request.status = new_status
        request.updated_at = datetime.utcnow()

        # 记录事件
        event = ApprovalEvent(
            event_id=f"evt-{len(request.history)}",
            action=action,
            actor=actor,
            comment=comment,
        )
        request.history.append(event)

        return new_status

    def get_available_actions(self, status: ApprovalStatus) -> list[ApprovalAction]:
        return [
            action for (s, action) in self._transitions
            if s == status
        ]

#3. 审批链构建器

#3.1 基于规则的审批链

Python
class ApprovalChainBuilder:
    """审批链构建器"""

    def __init__(self):
        self._rules: list[dict] = []

    def add_rule(self, condition: str,
                 levels: list[ApprovalLevel]) -> ApprovalChainBuilder:
        self._rules.append({
            "condition": condition,
            "levels": levels,
        })
        return self

    def build(self, request: ApprovalRequest) -> list[ApprovalLevel]:
        """根据请求内容构建审批链"""
        for rule in self._rules:
            if self._evaluate_condition(rule["condition"], request):
                return rule["levels"]

        # 默认单级审批
        return [ApprovalLevel(level=0, approvers=["default_approver"])]

    def _evaluate_condition(self, condition: str,
                            request: ApprovalRequest) -> bool:
        context = {
            "amount": request.amount,
            "domain": request.domain,
            "priority": request.priority,
        }
        try:
            return bool(eval(condition, {"__builtins__": {}}, context))
        except Exception:
            return False


# 使用示例
chain_builder = (
    ApprovalChainBuilder()
    .add_rule(
        "amount > 500000",
        [
            ApprovalLevel(0, ["department_manager"], mode="any", timeout_hours=24),
            ApprovalLevel(1, ["vp_finance", "vp_ops"], mode="all", timeout_hours=48),
            ApprovalLevel(2, ["ceo"], mode="any", timeout_hours=72),
        ]
    )
    .add_rule(
        "amount > 100000",
        [
            ApprovalLevel(0, ["department_manager"], mode="any", timeout_hours=24),
            ApprovalLevel(1, ["vp_finance"], mode="any", timeout_hours=48),
        ]
    )
    .add_rule(
        "amount > 0",
        [
            ApprovalLevel(0, ["department_manager"], mode="any", timeout_hours=48),
        ]
    )
)

#4. Temporal Workflow 实现

#4.1 审批 Workflow

Python
from temporalio import workflow, activity
from temporalio.common import RetryPolicy
from datetime import timedelta


@workflow.defn
class ApprovalWorkflow:
    """审批工作流"""

    def __init__(self):
        self._request: ApprovalRequest | None = None
        self._state_machine = ApprovalStateMachine()
        self._pending_action: ApprovalAction | None = None
        self._pending_actor: str = ""
        self._pending_comment: str = ""

    @workflow.run
    async def run(self, request_data: dict) -> dict:
        """主工作流"""
        self._request = ApprovalRequest(**request_data)

        # 构建审批链
        chain = await workflow.execute_activity(
            build_approval_chain,
            args=[request_data],
            start_to_close_timeout=timedelta(seconds=30),
        )
        self._request.approval_chain = [
            ApprovalLevel(**level) for level in chain
        ]

        # 提交审批
        self._state_machine.transition(
            self._request, ApprovalAction.SUBMIT,
            self._request.applicant
        )

        # 逐级审批
        for level_idx, level in enumerate(self._request.approval_chain):
            self._request.current_level = level_idx

            # 通知审批人
            await workflow.execute_activity(
                notify_approvers,
                args=[self._request.request_id,
                      level.approvers, level.level],
                start_to_close_timeout=timedelta(seconds=30),
            )

            # 等待审批(支持超时)
            result = await self._wait_for_approval(level)

            if result == "rejected":
                return self._build_result("rejected")
            elif result == "escalated":
                continue  # 升级到下一级
            elif result == "withdrawn":
                return self._build_result("withdrawn")

        # 所有层级通过
        self._state_machine.transition(
            self._request, ApprovalAction.APPROVE,
            "system", "All levels approved"
        )

        # 执行决策
        await workflow.execute_activity(
            execute_approved_decision,
            args=[self._request.decision_id],
            start_to_close_timeout=timedelta(seconds=60),
            retry_policy=RetryPolicy(maximum_attempts=3),
        )

        return self._build_result("approved")

    async def _wait_for_approval(self, level: ApprovalLevel) -> str:
        """等待审批人操作"""
        if level.mode == "any":
            return await self._wait_any(level)
        else:
            return await self._wait_all(level)

    async def _wait_any(self, level: ApprovalLevel) -> str:
        """或签:任意一人通过即可"""
        timeout = timedelta(hours=level.timeout_hours)

        try:
            # 等待 Signal
            await workflow.wait_condition(
                lambda: self._pending_action is not None,
                timeout=timeout,
            )

            action = self._pending_action
            self._pending_action = None

            if action == ApprovalAction.APPROVE:
                self._state_machine.transition(
                    self._request, ApprovalAction.APPROVE,
                    self._pending_actor, self._pending_comment
                )
                return "approved"
            elif action == ApprovalAction.REJECT:
                self._state_machine.transition(
                    self._request, ApprovalAction.REJECT,
                    self._pending_actor, self._pending_comment
                )
                return "rejected"
            elif action == ApprovalAction.WITHDRAW:
                return "withdrawn"
            else:
                return "approved"

        except TimeoutError:
            # 超时处理
            if level.auto_escalate_to:
                self._state_machine.transition(
                    self._request, ApprovalAction.ESCALATE,
                    "system", f"Timeout after {level.timeout_hours}h"
                )
                return "escalated"
            else:
                return "rejected"

    async def _wait_all(self, level: ApprovalLevel) -> str:
        """会签:所有人都需要通过"""
        remaining = set(level.approvers)
        timeout = timedelta(hours=level.timeout_hours)

        while remaining:
            try:
                await workflow.wait_condition(
                    lambda: self._pending_action is not None,
                    timeout=timeout,
                )

                action = self._pending_action
                actor = self._pending_actor
                self._pending_action = None

                if action == ApprovalAction.REJECT:
                    self._state_machine.transition(
                        self._request, ApprovalAction.REJECT,
                        actor, self._pending_comment
                    )
                    return "rejected"

                if action == ApprovalAction.APPROVE and actor in remaining:
                    remaining.discard(actor)

            except TimeoutError:
                return "escalated"

        return "approved"

    @workflow.signal
    async def approval_action(self, action: str, actor: str,
                               comment: str = ""):
        """接收审批动作的 Signal"""
        self._pending_action = ApprovalAction(action)
        self._pending_actor = actor
        self._pending_comment = comment

    @workflow.query
    def get_status(self) -> dict:
        """查询审批状态"""
        if self._request is None:
            return {"status": "not_started"}
        return {
            "request_id": self._request.request_id,
            "status": self._request.status.value,
            "current_level": self._request.current_level,
            "total_levels": len(self._request.approval_chain),
            "history": [
                {
                    "action": e.action.value,
                    "actor": e.actor,
                    "comment": e.comment,
                    "timestamp": e.timestamp.isoformat(),
                }
                for e in self._request.history
            ],
        }

    def _build_result(self, outcome: str) -> dict:
        return {
            "request_id": self._request.request_id,
            "outcome": outcome,
            "history_count": len(self._request.history),
        }

#4.2 Activities

Python
@activity.defn
async def build_approval_chain(request_data: dict) -> list[dict]:
    """构建审批链"""
    request = ApprovalRequest(**request_data)
    chain = chain_builder.build(request)
    return [
        {
            "level": level.level,
            "approvers": level.approvers,
            "mode": level.mode,
            "timeout_hours": level.timeout_hours,
            "auto_escalate_to": level.auto_escalate_to,
        }
        for level in chain
    ]


@activity.defn
async def notify_approvers(request_id: str,
                            approvers: list[str],
                            level: int) -> None:
    """通知审批人"""
    for approver in approvers:
        await notification_service.send(
            channel="email",
            recipient=approver,
            template="approval_pending",
            variables={
                "request_id": request_id,
                "level": level,
                "action_url": f"/approvals/{request_id}",
            },
        )


@activity.defn
async def execute_approved_decision(decision_id: str) -> dict:
    """执行已审批的决策"""
    result = await action_engine.execute(decision_id)
    return {"decision_id": decision_id, "execution_status": result.status}


@activity.defn
async def send_timeout_reminder(request_id: str,
                                 approver: str,
                                 hours_remaining: int) -> None:
    """发送超时提醒"""
    await notification_service.send(
        channel="im",
        recipient=approver,
        template="approval_reminder",
        variables={
            "request_id": request_id,
            "hours_remaining": hours_remaining,
        },
    )

#5. 审批 API

#5.1 gRPC 服务

PROTOBUF
syntax = "proto3";
package onto.approval.v1;

service ApprovalService {
    rpc SubmitApproval(SubmitRequest) returns (SubmitResponse);
    rpc ApproveOrReject(ActionRequest) returns (ActionResponse);
    rpc GetApprovalStatus(StatusRequest) returns (StatusResponse);
    rpc ListPendingApprovals(ListRequest) returns (ListResponse);
    rpc WithdrawApproval(WithdrawRequest) returns (WithdrawResponse);
}

message SubmitRequest {
    string decision_id = 1;
    string applicant = 2;
    string domain = 3;
    string title = 4;
    map<string, string> content = 5;
    double amount = 6;
    string priority = 7;
}

message ActionRequest {
    string request_id = 1;
    string action = 2;          // approve, reject, return
    string actor = 3;
    string comment = 4;
}

message StatusResponse {
    string request_id = 1;
    string status = 2;
    int32 current_level = 3;
    int32 total_levels = 4;
    repeated HistoryEntry history = 5;
}

#5.2 服务实现

Python
from temporalio.client import Client


class ApprovalServiceImpl:
    """审批 gRPC 服务"""

    def __init__(self, temporal_client: Client):
        self._client = temporal_client

    async def SubmitApproval(self, request, context):
        """提交审批"""
        request_id = f"apr-{datetime.utcnow().strftime('%Y%m%d%H%M%S')}"

        handle = await self._client.start_workflow(
            ApprovalWorkflow.run,
            {
                "request_id": request_id,
                "decision_id": request.decision_id,
                "applicant": request.applicant,
                "domain": request.domain,
                "title": request.title,
                "content": dict(request.content),
                "amount": request.amount,
                "priority": request.priority,
            },
            id=f"approval-{request_id}",
            task_queue="approval-tasks",
        )

        return {"request_id": request_id, "workflow_id": handle.id}

    async def ApproveOrReject(self, request, context):
        """审批/拒绝"""
        handle = self._client.get_workflow_handle(
            f"approval-{request.request_id}"
        )
        await handle.signal(
            ApprovalWorkflow.approval_action,
            args=[request.action, request.actor, request.comment],
        )
        return {"status": "action_recorded"}

    async def GetApprovalStatus(self, request, context):
        """查询审批状态"""
        handle = self._client.get_workflow_handle(
            f"approval-{request.request_id}"
        )
        status = await handle.query(ApprovalWorkflow.get_status)
        return status

#6. 超时与升级策略

#6.1 定时器管理

Python
class TimeoutManager:
    """审批超时管理"""

    @staticmethod
    async def setup_reminders(workflow_context,
                               level: ApprovalLevel,
                               request_id: str) -> None:
        """设置定时提醒"""
        total_hours = level.timeout_hours

        # 50% 时提醒
        await asyncio.sleep(total_hours * 0.5 * 3600)
        for approver in level.approvers:
            await workflow.execute_activity(
                send_timeout_reminder,
                args=[request_id, approver, int(total_hours * 0.5)],
                start_to_close_timeout=timedelta(seconds=30),
            )

        # 80% 时再次提醒
        await asyncio.sleep(total_hours * 0.3 * 3600)
        for approver in level.approvers:
            await workflow.execute_activity(
                send_timeout_reminder,
                args=[request_id, approver, int(total_hours * 0.2)],
                start_to_close_timeout=timedelta(seconds=30),
            )

#6.2 升级策略

Code
审批超时升级策略:

  时间线                                          动作
  ├──── 0h ─── 提交审批 ──────────────────────── 通知审批人
  │
  ├──── 24h ── 50% 提醒 ──────────────────────── 发送提醒
  │
  ├──── 38h ── 80% 提醒 ──────────────────────── 紧急提醒
  │
  ├──── 48h ── 超时 ──────────────────────────── 自动升级
  │                                               │
  │              ┌────────────────────────────────┘
  │              ▼
  │         升级到上级审批人
  │         重置超时计时器
  │
  └──── 96h ── 二次超时 ─────────────────────── 升级到 VP

#7. 审批看板

Code
审批看板视图:

  待审批 (12)          审批中 (5)           已完成 (89)
  ┌──────────────┐    ┌──────────────┐    ┌──────────────┐
  │ APR-0421     │    │ APR-0415     │    │ APR-0410 ✓   │
  │ 采购 ¥85,000 │    │ 预算 ¥500K   │    │ 差旅 ¥12,000 │
  │ 等待: 张经理  │    │ L2/3 VP审核  │    │ 已通过       │
  │ 剩余: 36h    │    │ 剩余: 12h    │    │              │
  ├──────────────┤    ├──────────────┤    ├──────────────┤
  │ APR-0420     │    │ APR-0413     │    │ APR-0408 ✗   │
  │ 合同 ¥200K   │    │ 人员 ¥800K   │    │ 采购 ¥50,000 │
  │ 等待: 李总监  │    │ L2/3 CFO会签 │    │ 已拒绝       │
  │ 剩余: 22h    │    │ 剩余: 45h    │    │              │
  └──────────────┘    └──────────────┘    └──────────────┘

#8. 与 DecisionEngine 集成

Python
class DecisionApprovalIntegration:
    """决策引擎与审批引擎集成"""

    def __init__(self, decision_engine, approval_service):
        self._decision = decision_engine
        self._approval = approval_service

    async def evaluate_with_approval(self, context) -> dict:
        """评估决策并按需触发审批"""
        result = self._decision.evaluate(context)

        needs_approval = self._check_approval_needed(result, context)

        if not needs_approval:
            return {
                "decision": result.decision,
                "status": "auto_approved",
                "approval_required": False,
            }

        # 触发审批流程
        approval_result = await self._approval.SubmitApproval({
            "decision_id": context.context_id,
            "applicant": context.metadata.get("applicant", "system"),
            "domain": context.domain,
            "title": f"Decision approval: {context.context_id}",
            "content": context.inputs,
            "amount": context.inputs.get("amount", 0),
        })

        return {
            "decision": result.decision,
            "status": "pending_approval",
            "approval_required": True,
            "approval_id": approval_result["request_id"],
        }

    def _check_approval_needed(self, result, context) -> bool:
        """判断是否需要审批"""
        # 低置信度决策需要审批
        if result.confidence < 0.7:
            return True
        # 高金额需要审批
        if context.inputs.get("amount", 0) > 10000:
            return True
        # 高风险域需要审批
        if context.domain in ("credit", "compliance", "hr"):
            return True
        return False

#9. 性能与可靠性

#9.1 性能指标

指标
Workflow 启动延迟< 100ms
Signal 处理延迟< 50ms
Query 响应时间< 20ms
并发 Workflow 数100,000+
单 Workflow 最大持续时间30 天

#9.2 可靠性保障

Code
Temporal 可靠性机制:

  ┌─────────────────────────────────────┐
  │           Temporal Server            │
  │                                     │
  │  ┌───────────┐  ┌───────────────┐   │
  │  │ Workflow   │  │ Event History │   │
  │  │ Execution  │  │ (持久化)      │   │
  │  └───────────┘  └───────────────┘   │
  │                                     │
  │  特性:                              │
  │  - 自动重试失败的 Activity            │
  │  - Worker 崩溃后自动恢复              │
  │  - 完整的事件历史审计                  │
  │  - 支持版本化 Workflow 升级            │
  └─────────────────────────────────────┘

#10. 实战案例

Python
# 场景:信贷决策触发多级审批

# 1. 决策引擎评估
context = (
    DecisionContextBuilder("credit")
    .with_inputs(
        credit_score=620,
        amount=300000,
        debt_ratio=0.45,
        applicant="user-12345",
    )
    .build()
)

# 2. 评估并触发审批
integration = DecisionApprovalIntegration(decision_engine, approval_service)
result = await integration.evaluate_with_approval(context)
# {
#   "decision": "conditional_approve",
#   "status": "pending_approval",
#   "approval_required": True,
#   "approval_id": "apr-20260324103000",
# }

# 3. 审批人操作
await approval_service.ApproveOrReject({
    "request_id": "apr-20260324103000",
    "action": "approve",
    "actor": "department_manager",
    "comment": "客户信用记录良好,批准",
})

# 4. 查询状态
status = await approval_service.GetApprovalStatus({
    "request_id": "apr-20260324103000",
})
# {
#   "status": "pending",
#   "current_level": 1,
#   "total_levels": 2,
#   "history": [...]
# }

#Key Takeaways

  1. 状态机 + Temporal 组合提供了持久化、可恢复、可审计的审批流程
  2. 审批链构建器 支持基于金额、域、优先级的动态审批链
  3. 或签/会签 通过 Temporal Signal 机制实现灵活的审批模式
  4. 超时升级 自动提醒并升级未及时处理的审批
  5. Temporal Query 支持实时查询审批进度,无需额外状态存储
  6. 与 DecisionEngine 集成 根据决策置信度和业务规则自动触发审批
  7. Event History 提供完整的审计链,满足合规要求

#Next Article

下一篇 S5-11 决策追踪链:从输入到执行的端到端可追溯 将详解如何构建从原始数据到最终执行的完整追踪链路。

tags: #approval-workflow #temporal #state-machine #multi-level #countersign #escalation #coomia-dip