← 返回文章列表

Dify 插件开发实验(05):有状态与幂等——插件如何安全地保持状态和处理重复调用?

1. 业务场景

先讲一个我们实际遇到的场景。

客服工单 SaaS 有一个事件通道:第三方系统会实时推送工单事件(创建/更新/关闭),平台收到后要记录、更新工单状态。但网络是脆弱的——推送方没收到确认就重试,一条「工单已创建」的事件可能在几秒内被推两次、三次。如果每次都当成新事件处理,同一张工单就会出现两条重复记录,状态还会被旧事件覆盖回去。

我们第一次做这类通道时,第一反应也是「收到事件就处理,处理完就完事」。真正动手才发现——「可能重复到达」才是常态,接收方必须自己扛起幂等:重复事件处理两次,工单记录就脏了;处理状态不落盘,排查只能靠猜;两个相同事件并发到达,先后都判「不存在」然后都写入,幂等形同虚设。

这不是个例。任何「可能重复到达」的数据通道都是这个模式:支付回调、工单事件、消息推送、Webhook 通知——发送方为了可靠性必然重试,接收方就必须自己处理「同一事件只处理一次」。

2. 场景痛点

这个流程的痛点,在事件通道上体现得最直接:

本质上,事件通道的可靠性不在发送方,而在接收方——「可能重复到达」是常态,幂等与有状态是接收方必须自己扛起来的能力。

3. 方案:为什么是插件化的 KV + 幂等

选这个方案,我们实际对比过:

这篇文章我们就用它搭一个事件接收工具插件:event_ingest(幂等写入)+ event_status(状态查询),跑通「重复事件只处理一次、并发不双写、状态可查」的完整链路。

4. 整体架构

graph TD subgraph app["【验证应用】"] start["开始(event_id/event_type/payload)"] ingest["接收事件(event_ingest)"] status["查询事件状态(event_status)"] out["输出(result_ingest + result_status)"] end1["结束"] end start --> ingest --> status --> out --> end1 subgraph inner["【插件内部】event_ingest 幂等判重"] kv["KV 查 event_id"] dup{"已存在?"} dupR["返回 {duplicate: true, status}(不重复处理)"] write["KV 写入 processing 态"] acc["返回 {accepted: true}"] st["event_status:KV 按 event_id 读回完整记录(有状态)"] end kv --> dup dup -- "是" --> dupR dup -- "否" --> write --> acc

链路很清晰:收事件 → 按 event_id 判重 → 首次写入 processing 态 → 查询读回完整记录。关键设计是「先查后写」的判重语义——重复事件直接返回当前状态,不覆盖不重入;KV 不可达时明确报错,绝不假装成功。

5. 模块设计

5.1 工具参数声明(tools/event_ingest.yaml)

event_id 是幂等键,来源方生成天然唯一:

parameters:

  - name: event_id

    type: string

    required: true

    form: llm

    label:

      zh_Hans: 事件 ID

    llm_description: 'Unique event id from the source system, e.g. EVT-20260805-001'

  - name: event_type

    type: string

    required: true

    form: llm

    llm_description: 'Event type, one of created/updated/closed'

  - name: payload

    type: string

    required: false

    form: llm

5.2 幂等判重核心逻辑(tools/event_ingest.py)

先查后写,判重与写入尽量原子:

# 幂等判重(KV 无原子 set-if-absent——极小竞态窗口,已记录)

existing, err_msg = kv_get(kv_url, key)

if err_msg:

    yield self.create_text_message(err_msg)

    return

if existing:

    yield self.create_text_message(json.dumps({

        "duplicate": True,

        "event_id": event_id,

        "status": existing.get("status", "unknown"),

        "received_at": existing.get("received_at"),

    }, ensure_ascii=False))

    return

# 首次接收:写入 processing 态

record = {"event_id": event_id, "event_type": event_type,

          "payload": payload, "status": "processing",

          "received_at": datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC")}

ok, err_msg = kv_set(kv_url, key, record)

if err_msg:

    yield self.create_text_message(err_msg)

    return

yield self.create_text_message(json.dumps({"accepted": True, "event_id": event_id,

    "status": record["status"], "received_at": record["received_at"]}, ensure_ascii=False))

5.3 KV 地址走凭证(provider/stateful_tool.yaml)

kv_url 默认 http://172.19.0.50:8123,换环境只改凭证不改代码。公共模块 tools/common.py 收敛常量与 KV 请求:KEY_PREFIX = "evt_"、错误分层 param_invalid / upstream_error / not_found——KV 不可达返回 upstream_error,绝不假装成功。

6. 运行验证

输入 预期 结果
首次 ingest 事件 EVT-TEST-001 accepted=true,status=processing ✅ 一致
同一 event_id 再次 ingest duplicate=true + 当前状态,不重复处理 ✅ 一致
ingest 后 event_status 查询 完整记录可读(有状态) ✅ 一致
两线程并发提交相同事件 只处理一次 ⚠️ 2 accepted / 0 duplicate(竞态窗口实测证实,预期内)
KV 不可达(错地址) 明确报错不假装成功 ✅ upstream_error
workflow 集成两轮冒烟 首次/重复幂等逻辑生效 ✅ 通过

环境:Dify 1.16.1(Docker Compose,daemon 0.6.1-local),KV 容器 dify104-kv(172.19.0.50:8123)。插件(daemon 容器)→ KV 直连可达,不经 ssrf_proxy(与工作流 http 节点不同,实测确认)。

7. 实战坑

现象 修复
KV key 含冒号 /state/evt:XXX → KV 400(expected str, bytes or os.PathLike object,路径解析问题) 前缀用下划线:evt_
幂等非原子 先查后写竞态,并发实测 2 accepted 接受并记录边界;生产用 Redis SETNX/唯一约束消除
常量重复定义 KEY_PREFIX 在 common.py + event_ingest.py 双份,旧值覆盖新值(冒号 key 根因) 常量单处维护,收敛到 common.py
KV 不可达假装成功 写入失败仍返回 accepted 会丢数据 失败返回 upstream_error 明确报错(迁移 105 静默失败教训)
插件出口白名单 误以为插件与 http 节点同受 ssrf_proxy 限制 实测 daemon 直连 docker 网段 IP 可达,不经代理无限制

8. 实验文档及源码获取

文章聚焦核心配置与采坑点,完整分步操作与验证记录见实验文档原文。

联系我

15088711270

手机端点击号码可直接拨打 · 桌面端可复制

微信二维码

扫码加微信 · 备注「门户」更快通过