Chen Yuan

我的 AI 代理只有一个任务:写一篇文章,把它发布到平台上,然后清理现场。但十次里有三次,它会在流水线中途崩溃,留下孤立的草稿,丢失临时文件,然后再次尝试发布同一篇文章。结果就是:重复的草稿、半失败的帖子,以及——最糟糕的情况——在同一信息流上出现两篇完全相同的文章。

我需要一种让流水线具备崩溃安全性的方法。不是“最终一致”的那种崩溃安全性,而是原子性的崩溃安全性。以下是实现它的状态机模式。

解决方案

我把流水线拆成多个阶段,每个阶段只负责一个职责。run ID 把所有内容串联起来。清单文件记录哪些临时文件属于哪个 run。状态机只有在当前阶段的所有门禁都通过时才会继续前进。

from __future__ import annotations
import json
from pathlib import Path
from enum import Enum
from datetime import datetime
from typing import Optional

class Phase(Enum):
    LOAD = "loading"
    SELECT = "topic_selection"
    DEDUP = "topic_dedup"
    WRITE = "article_generation"
    GATE = "de_ai_gate"
    VALIDATE = "pre_publish_validate"
    PUBLISH = "publish"
    VERIFY = "post_publish_verify"
    CLEAN = "cleanup"

class RunState:
    def __init__(self, run_id: str, manifest_path: Path):
        self.run_id = run_id
        self.manifest_path = manifest_path
        self.current_phase = None
        self.phase_results = {}

    def save(self):
        data = {
            "run_id": self.run_id,
            "current_phase": self.current_phase.value if self.current_phase else None,
            "phase_results": self.phase_results,
            "updated_at": datetime.now().isoformat(),
        }
        self.manifest_path.write_text(json.dumps(data, ensure_ascii=False, indent=2))

    @classmethod
    def load(cls, manifest_path: Path):
        if not manifest_path.exists():
            return None
        data = json.loads(manifest_path.read_text(encoding="utf-8"))
        state = cls(data["run_id"], manifest_path)
        state.current_phase = Phase(data["current_phase"]) if data["current_phase"] else None
        state.phase_results = data["phase_results"]
        return state

Enter fullscreen mode Exit fullscreen mode

状态机本身是一个简单的循环。每个阶段都是一个函数,返回 True(前进)或抛出异常(停止)。在前进之前,状态会被保存到磁盘。

class Pipeline:
    def __init__(self, state: RunState, temp_files: list[Path]):
        self.state = state
        self.temp_files = temp_files
        self.guard = Guard(temp_files)

    def run(self, phases: list[Phase]):
        for phase in phases:
            print(f"[{self.state.run_id}] Phase: {phase.value}")
            self.state.current_phase = phase
            self.state.save()

            handler = getattr(self, f"_phase_{phase.value}")
            handler()

            self.state.phase_results[phase.value] = True
            self.state.save()

    def _phase_cleanup(self):
        for f in self.temp_files:
            if f.exists():
                f.unlink()
                print(f"  cleaned: {f.name}")
        self.state.manifest_path.unlink(missing_ok=True)

Enter fullscreen mode Exit fullscreen mode

Guard 类位于流水线和文件系统之间。它准确知道哪些文件是临时发布产物,并在技能一致性检查期间忽略它们。

class Guard:
    def __init__(self, allowed_temps: list[Path]):
        self.allowed = {p.resolve() for p in allowed_temps}

    def is_temp_artifact(self, path: Path) -> bool:
        return path.resolve() in self.allowed

    def verify_no_untracked_writes(self, workspace: Path):
        for f in workspace.iterdir():
            if f.is_file() and not self.is_temp_artifact(f):
                raise RuntimeError(
                    f"Untracked file: {f.name}. "
                    "Add to temp manifest or refactor."
                )

Enter fullscreen mode Exit fullscreen mode

为什么有效

三个设计决策让这个模式具备了崩溃安全性。

第一,磁盘优先的状态。状态机在每个阶段处理程序执行前和执行后都会写入磁盘。如果代理在 _phase_validate() 期间崩溃,下次运行会读取 manifest.json,看到 current_phase: "pre_publish_validate",从而准确知道从哪里恢复。不需要猜测“上一次做到哪了”。

第二,临时文件清单。在流水线启动前,它将要创建的每个文件都会被注册。Guard 类会强制执行这一点:如果写入的文件不在清单中,流水线就会停止。这让我抓到过两次——一次是我在 Python 中使用了 write_file 而不是 open(),另一次是日志库在工作目录创建了 .log 文件。

第三,幂等的阶段处理程序。每个阶段函数在再次执行前会检查工作是否已完成。_phase_dedup() 会读取现有的文章列表,只有在尚未写入的情况下才会写入新数据。这意味着在崩溃后重放第 2 阶段是安全的——它不会重复检查或重复写入。

def _phase_dedup(self):
    manifest = self.state.manifest_path
    dedup_flag = manifest.parent / f"{self.state.run_id}_dedup_done"
    if dedup_flag.exists():
        print("  dedup already done, skipping")
        return
    # ... actual dedup logic ...
    dedup_flag.write_text("done")

Enter fullscreen mode Exit fullscreen mode

注意事项

Windows 的文件锁定机制不同。在 Linux 上,fcntl.flock() 能正常工作。在 Windows 上,你需要 msvcrt.locking()。当两个流水线运行意外重叠并试图写入同一个清单文件时,我遇到了这个问题。解决方案是跨平台锁:

import sys

def safe_write(path: Path, data: str):
    if sys.platform == "win32":
        import msvcrt
        with open(path, "w", encoding="utf-8") as f:
            msvcrt.locking(f.fileno(), msvcrt.LK_LOCK, 1)
            f.write(data)
            msvcrt.locking(f.fileno(), msvcrt.LK_UNLCK, 1)
    else:
        import fcntl
        with open(path, "w", encoding="utf-8") as f:
            fcntl.flock(f.fileno(), fcntl.LOCK_EX)
            f.write(data)
            fcntl.flock(f.fileno(), fcntl.LOCK_UN)

Enter fullscreen mode Exit fullscreen mode

清单文件本身也是临时文件。如果你把它包含在临时文件清单中,Guard 会允许它通过。但 _phase_cleanup 处理程序会最后删除清单——在所有其他临时文件都被清理后——这样在清理过程中崩溃也不会丢失状态记录。

状态机的效果取决于阶段边界的划分。我最初把阶段划分得过于细致(10 步流水线用了 14 个阶段)。状态文件被频繁覆盖,I/O 开销每轮增加了约 2 秒。我把相关步骤合并了:de_ai_gate + validate 合并为一个阶段,publish + verify 合并为另一个阶段。最终降到 7 个阶段,没有明显开销。

你怎么看?

我见过有人用 SQLite 保存流水线状态,用 Redis 处理分布式流水线,甚至在 S3 上用 JSON 文件。我的方法刻意保持最小化——一个 JSON 文件、一个 Guard 类和一个循环。它之所以有效,是因为流水线是单线程且运行在一台机器上。

你如何实现崩溃安全的自动化?你会选择工作流引擎(Temporal、Prefect)还是自己实现状态机?我很想听听你的经验。