Python 基础体系 · 第 104/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 自动化脚本:文件、命令、网络、重试、幂等和审计

自动化脚本最初往往只有几行:

download(url)
run_command()
move_file()

真正进入生产环境后,问题会迅速变成:

  • 下载到一半时进程被杀死,目标文件是否会留下半份内容?
  • 外部命令返回非零退出码时,是否应该重试?
  • 网络请求超时,但服务器其实已经完成了写入,重试会不会重复创建数据?
  • 脚本运行到一半崩溃,再执行一次是否安全?
  • 日志能否回答“谁在什么时间对哪个对象做了什么,结果如何”?
  • 两个定时任务同时启动时,是否会互相覆盖或重复执行?

因此,自动化脚本的核心不是“把几个 API 串起来”,而是设计一个能够面对部分完成、未知结果、重复执行和可追溯性的状态转换系统。

本文以 Python 3.14 标准库为范围,使用 pathlibsubprocessurllib.requesttempfileoshashlibloggingsqlite3argparse 构建一个可运行的自动化任务。Python 标准库将文件系统、子进程、网络、日志、命令行和 SQLite 分别提供为独立模块;自动化脚本需要做的是明确它们之间的数据流和失败边界。(docs.python.org)


一、先建立模型:自动化任务是一个有副作用的状态机

设一个任务从外部世界读取输入,执行副作用,然后写入结果:

S0读取S1执行S2提交S3S_0 \xrightarrow{\text{读取}} S_1 \xrightarrow{\text{执行}} S_2 \xrightarrow{\text{提交}} S_3

其中:

  • S0S_0:尚未开始;
  • S1S_1:输入已经确定,例如 URL、目标路径和命令参数;
  • S2S_2:副作用正在发生,例如下载文件、执行命令;
  • S3S_3:结果已提交并可被后续任务观察。

难点在于,进程可能在任意时刻崩溃。尤其是 S2S_2S3S_3 之间,调用方可能无法判断副作用究竟有没有完成。

例如:

客户端发送“创建订单”
服务器创建成功
服务器准备返回响应
网络连接断开
客户端收到超时

客户端看到的是失败,但服务器状态可能已经成功。此时直接重试并不能保证安全。

因此,一个可靠的自动化任务必须区分三类结果:

结果 含义 是否适合直接重试
成功 已知副作用完成 通常不重试
明确失败 已知副作用未完成,或操作明确拒绝 根据错误类型决定
未知 无法判断副作用是否完成 先查询、校验或使用幂等键

这也是“重试”和“幂等”必须放在一起讨论的原因:重试解决暂时性失败,幂等解决重复执行造成的业务风险


二、文件自动化:路径、编码、临时文件和原子提交

2.1 用 Path 表示路径

pathlib 将路径表示为对象,并区分只做路径计算的纯路径和能够访问文件系统的具体路径;通常使用当前平台对应的 Path。(docs.python.org)

from pathlib import Path

root = Path("workspace")
source = root / "input" / "data.txt"
target = root / "output" / "data.txt"

print(source)
print(source.parent)
print(source.suffix)

这段代码只进行路径拼接和属性读取,并不会自动创建目录或文件。路径对象的运算不等于文件系统操作:

target.parent.mkdir(parents=True, exist_ok=True)

这里才会创建目录。parents=True 表示缺失的父目录一并创建,exist_ok=True 表示目录已经存在时不视为错误。

常见误区是把字符串拼接当作路径拼接:

# 不推荐
path = base + "/" + filename

这会把路径分隔符、Windows 路径和特殊字符处理交给调用者。更重要的是,路径拼接不会自动防止路径穿越:

user_name = "../outside.txt"
target = root / user_name

Path 会得到一个合法路径对象,但它可能指向 root 之外。若文件名来自用户或网络数据,应显式校验解析后的路径:

def safe_child(root: Path, name: str) -> Path:
    root = root.resolve()
    candidate = (root / name).resolve()

    try:
        candidate.relative_to(root)
    except ValueError as exc:
        raise ValueError(f"非法路径: {name!r}") from exc

    return candidate

这里的因果关系是:

  1. / 只完成路径组合;
  2. resolve()... 和符号链接解析为实际路径;
  3. relative_to(root) 能验证候选路径是否仍位于根目录下;
  4. 只有通过验证后,才允许进行写操作。

这仍不是完整的安全边界。如果攻击者能够在校验后、打开文件前替换符号链接,就可能形成 TOCTOU(time-of-check to time-of-use)竞态。高权限程序应进一步使用更底层的目录文件描述符、受限权限和平台相关的安全打开方式。


2.2 文本和二进制必须分开处理

文本文件涉及:

  • 字节到字符的编码;
  • 换行符转换;
  • 解码失败策略;
  • 文件是否以追加还是覆盖方式打开。

二进制文件则不应该经过文本编码层:

from pathlib import Path

text_path = Path("message.txt")
text_path.write_text("你好\n", encoding="utf-8")

content = text_path.read_text(encoding="utf-8")
print(content)

对于网络下载、压缩包、图片、签名文件等内容,应使用二进制模式:

data = Path("archive.bin").read_bytes()
Path("archive-copy.bin").write_bytes(data)

不要用默认编码读取外部文件:

# 风险:默认编码取决于环境
text = Path("message.txt").read_text()

同一份脚本在本地终端、Linux 定时任务和 Windows 主机上可能得到不同的默认编码行为。自动化任务应把编码写在接口边界上,而不是依赖执行环境。


2.3 直接写目标文件会产生半成品

下面的代码存在明显故障窗口:

target = Path("result.json")

with target.open("w", encoding="utf-8") as f:
    f.write(large_content)

执行过程如下:

旧 result.json 存在
        |
打开 result.json,通常会先截断旧内容
        |
写入一部分
        |
进程崩溃
        |
result.json 只剩半份内容

如果其他程序正在读取 result.json,它可能看到:

  • 不完整的 JSON;
  • 截断的 CSV;
  • 长度正确但内容不完整的文件;
  • 空文件。

更安全的做法是:

  1. 在目标目录创建临时文件;
  2. 将完整内容写入临时文件;
  3. 刷新 Python 缓冲区;
  4. 使用 os.fsync() 请求将文件内容刷新到磁盘;
  5. os.replace() 替换目标文件。

os.fsync() 用于强制刷新文件描述符对应的文件;对 Python 文件对象,应先 flush(),再对其文件描述符调用 os.fsync()os.replace() 用于以替换语义重命名路径。(docs.python.org)

from __future__ import annotations

import os
import tempfile
from pathlib import Path


def atomic_write_text(
    target: Path,
    content: str,
    *,
    encoding: str = "utf-8",
) -> None:
    target.parent.mkdir(parents=True, exist_ok=True)

    fd, temporary_name = tempfile.mkstemp(
        prefix=f".{target.name}.",
        suffix=".tmp",
        dir=target.parent,
        text=True,
    )

    temporary_path = Path(temporary_name)

    try:
        with os.fdopen(fd, "w", encoding=encoding, newline="") as f:
            f.write(content)
            f.flush()
            os.fsync(f.fileno())

        os.replace(temporary_path, target)
    except BaseException:
        try:
            temporary_path.unlink()
        except FileNotFoundError:
            pass
        raise

为什么临时文件必须和目标文件位于同一个目录?

因为很多文件系统只能对同一文件系统内的重命名提供原子替换语义。如果临时文件放在系统临时目录,而目标文件位于另一个挂载点,代码可能退化为复制,或者直接失败。

这个函数保证的是:

观察者看到的目标状态{旧完整文件,新完整文件}\text{观察者看到的目标状态} \in \{\text{旧完整文件}, \text{新完整文件}\}

它不保证:

  • 目标文件所在目录已经持久化到断电后的存储介质;
  • 网络文件系统具有与本地文件系统相同的原子语义;
  • 文件内容和审计记录同时提交;
  • 多个进程不会同时产生不同版本并互相覆盖。

如果需要更强的崩溃持久性,通常还需要在替换后打开父目录并同步目录项,但这涉及操作系统差异,不能简单宣称在所有平台上具有相同语义。


三、命令自动化:参数边界、退出码、输出和超时

3.1 使用参数数组,而不是拼接 shell 字符串

subprocess.run() 会等待命令完成,并返回 CompletedProcess;它支持捕获输出、设置环境、工作目录、超时和 check 等参数。默认情况下,Python 不会隐式调用系统 shell。(docs.python.org)

推荐:

import subprocess

result = subprocess.run(
    ["python", "--version"],
    capture_output=True,
    text=True,
    check=False,
)

print("退出码:", result.returncode)
print("标准输出:", result.stdout)
print("标准错误:", result.stderr)

不推荐:

import os

filename = "report; rm -rf data"
os.system(f"processor {filename}")

或者:

subprocess.run(
    f"processor {filename}",
    shell=True,
)

shell=True 且命令字符串包含外部输入时,空格、引号和 shell 元字符都会影响命令解析,可能形成 shell 注入。官方文档明确要求调用方自行保证引用和转义正确。(docs.python.org)

参数数组的意义是保留参数边界:

filename = "report; rm -rf data"

subprocess.run(
    ["processor", filename],
    check=True,
)

这里的 filename 是一个参数值,而不是一段需要重新解析的 shell 语法。


3.2 退出码不是布尔值

进程退出码通常约定:

  • 0:成功;
  • 非零:失败或特殊状态;
  • Unix 系统中,负值可能表示被信号终止;
  • Windows 和具体命令可能有自己的约定。

因此不要只写:

if result.returncode:
    print("失败")

应结合具体命令的协议解释退出码:

def run_checked(command: list[str]) -> str:
    result = subprocess.run(
        command,
        capture_output=True,
        text=True,
        timeout=30,
        check=False,
    )

    if result.returncode != 0:
        raise RuntimeError(
            f"命令失败: command={command!r}, "
            f"returncode={result.returncode}, "
            f"stderr={result.stderr[-2000:]!r}"
        )

    return result.stdout

check=True 会在非零退出码时抛出 CalledProcessError,适合“非零就是异常”的简单场景:

subprocess.run(
    ["python", "-c", "raise SystemExit(2)"],
    check=True,
)

但在需要区分“目标不存在”“没有变化”“部分成功”等业务状态时,应该保留 returncode 并由应用层解释。


3.3 超时只解决等待问题,不自动撤销副作用

subprocess.run(timeout=...) 的超时最终通过 Popen.communicate() 处理;等待超时后,子进程会被杀死并等待,随后重新抛出 TimeoutExpired。进程创建本身在许多平台上不能被该超时中断。(docs.python.org)

import subprocess

try:
    result = subprocess.run(
        ["some-command"],
        capture_output=True,
        text=True,
        timeout=10,
        check=False,
    )
except subprocess.TimeoutExpired as exc:
    print("命令等待超时:", exc)

重要边界是:杀死子进程不等于撤销子进程已经完成的操作

例如:

命令已经上传 800 MB
命令继续等待远端确认
10 秒超时
Python 杀死命令

远端可能已经拥有完整文件。此时再次执行上传命令,可能产生重复对象或覆盖对象。

若需要手动使用 Popen,超时后的清理顺序应包括杀死子进程并再次 communicate(),以完成管道读取和回收;官方文档特别说明,超时后应通过额外的 communicate() 完成处理,而不是直接调用 wait()。(docs.python.org)

import subprocess


def run_process(command: list[str], timeout: float) -> subprocess.CompletedProcess[str]:
    process = subprocess.Popen(
        command,
        stdout=subprocess.PIPE,
        stderr=subprocess.PIPE,
        text=True,
    )

    try:
        stdout, stderr = process.communicate(timeout=timeout)
    except subprocess.TimeoutExpired:
        process.kill()
        stdout, stderr = process.communicate()
        raise RuntimeError(
            f"命令超时: command={command!r}, "
            f"stdout={stdout[-1000:]!r}, stderr={stderr[-1000:]!r}"
        )

    return subprocess.CompletedProcess(
        command,
        process.returncode,
        stdout,
        stderr,
    )

3.4 管道应优先显式连接

shell 中的:

producer | consumer

不等于启动一个命令。它至少包含:

  1. 启动 producer
  2. 创建标准输出管道;
  3. 启动 consumer
  4. 将管道连接到 consumer 的标准输入;
  5. 同时读取输出并等待两个进程;
  6. 处理其中一个进程提前退出的情况。

Python 可以显式表达:

import subprocess

producer = subprocess.Popen(
    ["printf", "alpha\nbeta\n"],
    stdout=subprocess.PIPE,
    text=True,
)

consumer = subprocess.Popen(
    ["grep", "beta"],
    stdin=producer.stdout,
    stdout=subprocess.PIPE,
    stderr=subprocess.PIPE,
    text=True,
)

assert producer.stdout is not None
producer.stdout.close()

stdout, stderr = consumer.communicate()
producer_returncode = producer.wait()

if producer_returncode != 0:
    raise RuntimeError(f"producer failed: {producer_returncode}")

if consumer.returncode != 0:
    raise RuntimeError(f"consumer failed: {consumer.returncode}")

print(stdout, end="")

关闭 producer.stdout 的父进程副本很重要,否则上游进程可能收不到下游关闭管道后产生的 SIGPIPE。官方 HOWTO 也采用了这一处理方式。(docs.python.org)


四、网络自动化:超时、状态码、响应体和未知结果

4.1 网络请求至少要有超时

from urllib.request import Request, urlopen

request = Request(
    "https://example.com/data.json",
    headers={"User-Agent": "wr-automation/1.0"},
)

with urlopen(request, timeout=15) as response:
    body = response.read()

urlopen()timeout 以秒为单位,用于连接等阻塞操作;文档说明该超时能力适用于 HTTP、HTTPS 和 FTP。(docs.python.org)

没有超时的网络调用可能让定时任务永久占用执行槽位:

调度器启动任务
任务连接远端
远端不发送数据,也不关闭连接
任务一直阻塞
下一轮任务启动
多个任务逐渐堆积

连接超时、读取超时、DNS 失败和 HTTP 错误不是同一种错误。urllib.request 使用 HTTPError 等异常表示 HTTP 层错误;重定向也可能需要客户端处理。(docs.python.org)

一个用于下载的实现可以这样写:

from __future__ import annotations

import hashlib
from pathlib import Path
from urllib.error import HTTPError, URLError
from urllib.request import Request, urlopen


class DownloadError(RuntimeError):
    pass


def download_bytes(url: str, *, timeout: float = 15) -> tuple[bytes, str]:
    request = Request(
        url,
        headers={"User-Agent": "wr-automation/1.0"},
    )

    try:
        with urlopen(request, timeout=timeout) as response:
            status = getattr(response, "status", None)
            if status is not None and not 200 <= status < 300:
                raise DownloadError(f"非成功状态码: {status}")

            body = response.read()
    except HTTPError as exc:
        raise DownloadError(
            f"HTTP 错误: status={exc.code}, url={url!r}"
        ) from exc
    except URLError as exc:
        raise DownloadError(f"网络错误: url={url!r}, reason={exc.reason!r}") from exc

    digest = hashlib.sha256(body).hexdigest()
    return body, digest


body, digest = download_bytes("https://example.com/data.json")
print(len(body), digest)

这里使用 SHA-256 的目的不是加密,而是生成内容指纹:

d=H(body)d = H(\text{body})

其中 HH 是哈希函数,d 用于判断两次得到的字节内容是否相同。哈希不能证明内容来自可信来源;如果需要来源认证,应使用签名、HMAC 或可信传输链路。


4.2 不同 HTTP 状态码应有不同策略

重试不能写成:

for _ in range(5):
    try:
        return request()
    except Exception:
        pass

至少应先区分:

状态或异常 通常含义 一般策略
2xx 请求成功 提交结果
400、401、403 请求或权限问题 通常不重试
404 资源不存在 通常不重试,除非资源有延迟发布语义
408、429 请求超时或被限流 可重试,优先遵守服务端提示
500、502、503、504 服务端或网关暂时失败 通常可重试
DNS、连接重置、连接超时 网络路径失败 通常可重试
读取超时 可能失败,也可能服务端已完成 对写操作必须考虑未知结果

对于 429,如果响应包含 Retry-After,应优先使用服务端给出的等待时间。否则,客户端只能使用有限的指数退避。


五、重试:只重试暂时性失败,并控制总预算

5.1 重试的数学模型

设第 kk 次重试前的基础等待时间为:

bk=min(C,B2k)b_k = \min(C, B \cdot 2^k)

其中:

  • BB:初始退避时间;
  • CC:最大退避上限;
  • k=0,1,2,k=0,1,2,\ldots:已经失败的次数。

为了避免多个客户端同时重试形成“惊群”,加入抖动:

wk=min(C,B2k)+U(0,J)w_k = \min(C, B \cdot 2^k) + U(0,J)

其中 U(0,J)U(0,J)00JJ 之间的随机等待。

例如 B=1B=1C=8C=8J=0.5J=0.5,基础等待序列是:

第 1 次失败后:1 秒 + 随机抖动
第 2 次失败后:2 秒 + 随机抖动
第 3 次失败后:4 秒 + 随机抖动
第 4 次失败后:8 秒 + 随机抖动
第 5 次失败后:8 秒 + 随机抖动

但“重试 5 次”不等于总耗时最多 5 秒。还要计算:

Ttotal=Trequest,i+wiT_{\text{total}} = \sum T_{\text{request},i} + \sum w_i

所以生产任务通常需要同时限制:

  • 单次请求超时;
  • 最大尝试次数;
  • 总时间预算;
  • 单任务并发数;
  • 对服务端的请求速率。

5.2 一个可测试的重试函数

from __future__ import annotations

import random
import time
from collections.abc import Callable
from typing import TypeVar

T = TypeVar("T")


def retry(
    operation: Callable[[], T],
    *,
    is_retryable: Callable[[Exception], bool],
    max_attempts: int = 4,
    base_delay: float = 1.0,
    max_delay: float = 8.0,
    jitter: float = 0.2,
    sleep: Callable[[float], None] = time.sleep,
) -> T:
    if max_attempts < 1:
        raise ValueError("max_attempts 必须至少为 1")

    for attempt in range(1, max_attempts + 1):
        try:
            return operation()
        except Exception as exc:
            last_attempt = attempt == max_attempts

            if last_attempt or not is_retryable(exc):
                raise

            exponential = min(max_delay, base_delay * 2 ** (attempt - 1))
            delay = exponential + random.uniform(0, jitter)
            sleep(delay)

    raise AssertionError("循环逻辑不应到达这里")

测试时可以注入假的 sleep,避免测试真的等待:

def test_retry():
    calls = 0
    waits: list[float] = []

    class TemporaryError(Exception):
        pass

    def operation() -> str:
        nonlocal calls
        calls += 1
        if calls < 3:
            raise TemporaryError()
        return "ok"

    result = retry(
        operation,
        is_retryable=lambda exc: isinstance(exc, TemporaryError),
        sleep=waits.append,
        jitter=0,
    )

    assert result == "ok"
    assert calls == 3
    assert waits == [1.0, 2.0]

这个测试验证了三个中间状态:

第 1 次:失败,等待 1 秒
第 2 次:失败,等待 2 秒
第 3 次:成功,返回 ok

不要捕获 BaseException 进行重试,因为 KeyboardInterruptSystemExit 等控制流异常不应被普通网络重试吞掉。上面的通用函数捕获 Exception,但仍要求调用方准确实现 is_retryable()


5.3 重试与操作类型必须配套

重试安全性不能只由异常类型决定,还取决于操作语义。

读取操作

GET /report

若服务器没有副作用,超时后重试通常风险较低,但仍要考虑远端限流和响应体重复下载。

可重复覆盖操作

PUT /objects/report-2026.json

如果相同路径写入相同内容,重复执行通常可以设计为幂等。客户端可以用内容哈希验证结果。

创建操作

POST /orders

超时后重试可能创建两个订单。必须使用服务端支持的幂等键:

Idempotency-Key: 7e7c...

或改为:

PUT /orders/client-generated-id

不可逆副作用

发送短信
扣款
删除远端数据
执行数据库迁移

不能因为“连接超时”就自动重试。应先查询状态,或让服务端提供幂等协议。


六、幂等:把重复执行变成同一个最终状态

6.1 幂等的形式化定义

设操作为函数 ff,状态为 xx。如果:

f(f(x))=f(x)f(f(x)) = f(x)

则称该操作在状态空间上是幂等的。

例如:

target.write_text("固定内容", encoding="utf-8")

如果每次都把同样内容写入同一个目标文件,最终内容相同,这个业务操作可以视为幂等。

但下面的操作不是幂等的:

with log.open("a", encoding="utf-8") as f:
    f.write("event\n")

执行两次会产生两条记录。

更容易混淆的是“结果相同”和“操作没有重复副作用”不是一回事:

counter += 1

即使最终脚本输出相同,计数器已经被增加两次。


6.2 幂等键、去重表和结果缓存

对于外部副作用,可以定义业务键:

K=任务类型业务对象 ID输入版本K = \text{任务类型} \Vert \text{业务对象 ID} \Vert \text{输入版本}

例如:

publish:report:2026-09-01:sha256=abc123

执行前先查询该键:

不存在 -> 执行
存在且成功 -> 直接返回历史结果
存在且运行中 -> 等待、退出或接管
存在但失败 -> 根据策略重试或人工处理

SQLite 适合作为单机脚本的持久化状态表:

import sqlite3
from pathlib import Path


def open_state_db(path: Path) -> sqlite3.Connection:
    connection = sqlite3.connect(path)
    connection.execute("""
        CREATE TABLE IF NOT EXISTS task_runs (
            idempotency_key TEXT PRIMARY KEY,
            status TEXT NOT NULL,
            result_hash TEXT,
            started_at TEXT NOT NULL,
            finished_at TEXT,
            error TEXT
        )
    """)
    connection.commit()
    return connection

PRIMARY KEY 使相同幂等键无法插入两次。竞争条件下,不能只写:

if not exists(key):
    insert(key)
    perform_side_effect()

两个进程可能同时通过 not exists,然后都执行副作用。正确的顺序通常是:

  1. 事务中抢占幂等键;
  2. 只有插入成功的进程执行副作用;
  3. 完成后更新为 succeeded
  4. 崩溃后由后续任务根据状态和租约恢复。

但数据库状态和外部副作用仍然不是一个原子事务:

数据库写入 running
执行远端上传成功
进程在更新 succeeded 前崩溃

下一次执行看到 running,却不知道上传是否已经成功。这就是“未知结果”问题。解决方法包括:

  • 查询远端对象是否存在;
  • 根据内容哈希比较;
  • 使用外部系统的幂等键;
  • 将远端对象命名为确定性键;
  • 设计补偿操作;
  • 对长任务使用租约和心跳。

6.3 文件发布的幂等实现

下面的发布逻辑使用内容哈希判断是否需要替换:

import hashlib
from pathlib import Path


def sha256_bytes(data: bytes) -> str:
    return hashlib.sha256(data).hexdigest()


def publish_if_changed(target: Path, data: bytes) -> str:
    new_hash = sha256_bytes(data)

    if target.exists():
        old_hash = sha256_bytes(target.read_bytes())
        if old_hash == new_hash:
            return "unchanged"

    target.parent.mkdir(parents=True, exist_ok=True)

    temporary = target.with_name(f".{target.name}.{new_hash}.tmp")
    temporary.write_bytes(data)

    with temporary.open("rb") as f:
        import os
        f.flush()
        os.fsync(f.fileno())

    import os
    os.replace(temporary, target)
    return "updated"

这个例子体现了两个层次:

  • 业务幂等性:相同内容重复发布,最终文件内容不变;
  • 文件提交安全:写入临时文件后再替换,避免目标文件半成品。

不过它仍有并发问题:

进程 A 读取旧文件
进程 B 读取旧文件
进程 A 写入新版本 A
进程 B 写入新版本 B
进程 B 替换掉 A

如果发布顺序必须可控,应增加版本号、文件锁或数据库中的 compare-and-set 条件:

只有当前版本仍为 v1,才能替换为 v2

七、审计:日志记录“发生了什么”,审计记录“谁改变了什么”

7.1 日志和审计不是同一件事

日志主要用于诊断程序运行过程:

连接超时
重试第 2 次
命令退出码为 1

审计用于回答责任和状态问题:

主体 wr-bot
在 2026-09-01T10:00:00+08:00
使用任务 run-abc
对对象 report.json
执行 publish
输入哈希 sha256=...
结果 updated

日志可以丢失、滚动或被大量调试信息淹没。审计记录应具有稳定字段和明确的状态转换。

一个审计事件至少应包含:

字段 含义
event_id 事件唯一标识
run_id 一次任务运行的标识
actor 执行主体
action 操作名称
object 被操作对象
input_hash 输入内容或参数摘要
attempt 第几次尝试
status started、succeeded、failed、unknown
started_atfinished_at 时间
error 结构化错误信息

时间应使用带时区的 ISO 8601 时间,而不是裸的本地时间字符串:

from datetime import datetime, timezone

now = datetime.now(timezone.utc).isoformat()
print(now)

展示给杭州用户时可以转换为 Asia/Shanghai,但内部事件最好保留明确时区。否则夏令时、跨机器日志和排序都会产生歧义。


7.2 使用结构化日志

Python 的 logging 会为每条记录创建 LogRecord,记录名称、级别、路径、行号、消息和异常信息等上下文;可以通过过滤器或记录工厂注入额外字段。(docs.python.org)

简单的 JSON 日志格式:

from __future__ import annotations

import json
import logging
import sys
from datetime import datetime, timezone


class JsonFormatter(logging.Formatter):
    def format(self, record: logging.LogRecord) -> str:
        payload = {
            "timestamp": datetime.now(timezone.utc).isoformat(),
            "level": record.levelname,
            "logger": record.name,
            "message": record.getMessage(),
            "run_id": getattr(record, "run_id", None),
            "action": getattr(record, "action", None),
        }

        if record.exc_info:
            payload["exception"] = self.formatException(record.exc_info)

        return json.dumps(payload, ensure_ascii=False)


def configure_logging() -> logging.Logger:
    handler = logging.StreamHandler(sys.stderr)
    handler.setFormatter(JsonFormatter())

    logger = logging.getLogger("wr_automation")
    logger.setLevel(logging.INFO)
    logger.handlers.clear()
    logger.addHandler(handler)
    logger.propagate = False
    return logger

调用:

logger = configure_logging()

logger.info(
    "开始发布",
    extra={
        "run_id": "run-001",
        "action": "publish",
    },
)

输出类似:

{
  "timestamp": "2026-09-01T02:00:00+00:00",
  "level": "INFO",
  "logger": "wr_automation",
  "message": "开始发布",
  "run_id": "run-001",
  "action": "publish"
}

不要把密码、访问令牌、完整 Cookie 和未经处理的用户隐私直接写入日志。哈希也不是总能隐藏隐私:低熵数据可以被枚举反推。审计字段应记录“足够定位问题的信息”,而不是“所有可获得的信息”。


7.3 审计状态必须覆盖未知结果

一个只记录成功和失败的审计模型是不完整的:

started -> failed
started -> succeeded

至少还需要:

started -> unknown

例如网络写操作超时:

try:
    response = send_request()
except TimeoutError:
    audit(
        status="unknown",
        error="请求超时,无法判断远端是否已提交",
    )
    raise

unknown 不是“失败的另一种写法”,它表示系统缺少事实。后续恢复流程必须先查询或校验,而不是盲目重新执行。


八、端到端示例:下载、校验、执行命令、原子发布和审计

下面的示例实现一个简化任务:

  1. 接收 URL 和输出路径;
  2. 下载内容;
  3. 计算 SHA-256;
  4. 通过临时文件和 os.replace() 发布;
  5. 运行一个外部校验命令;
  6. 记录结构化日志;
  7. argparse 暴露命令行接口。

argparse 支持子命令;通过 add_subparsers() 可以把不同操作分成独立命令,并用 set_defaults() 绑定处理函数。(docs.python.org)

from __future__ import annotations

import argparse
import hashlib
import json
import logging
import os
import subprocess
import sys
import tempfile
import uuid
from datetime import datetime, timezone
from pathlib import Path
from urllib.error import HTTPError, URLError
from urllib.request import Request, urlopen


class AutomationError(RuntimeError):
    pass


class RetryableDownloadError(AutomationError):
    pass


def utc_now() -> str:
    return datetime.now(timezone.utc).isoformat()


def digest_bytes(data: bytes) -> str:
    return hashlib.sha256(data).hexdigest()


def download(url: str, timeout: float) -> bytes:
    request = Request(
        url,
        headers={"User-Agent": "wr-automation/1.0"},
    )

    try:
        with urlopen(request, timeout=timeout) as response:
            status = getattr(response, "status", 200)

            if status in {408, 429} or 500 <= status <= 599:
                raise RetryableDownloadError(f"可重试 HTTP 状态码: {status}")

            if not 200 <= status <= 299:
                raise AutomationError(f"不可重试 HTTP 状态码: {status}")

            return response.read()

    except HTTPError as exc:
        if exc.code == 429 or 500 <= exc.code <= 599:
            raise RetryableDownloadError(
                f"可重试 HTTP 错误: {exc.code}"
            ) from exc
        raise AutomationError(f"HTTP 错误: {exc.code}") from exc

    except URLError as exc:
        raise RetryableDownloadError(
            f"网络错误: {exc.reason!r}"
        ) from exc


def atomic_publish(target: Path, data: bytes) -> str:
    target.parent.mkdir(parents=True, exist_ok=True)

    new_hash = digest_bytes(data)

    if target.exists():
        old_hash = digest_bytes(target.read_bytes())
        if old_hash == new_hash:
            return "unchanged"

    fd, temporary_name = tempfile.mkstemp(
        prefix=f".{target.name}.",
        suffix=".tmp",
        dir=target.parent,
    )
    temporary = Path(temporary_name)

    try:
        with os.fdopen(fd, "wb") as f:
            f.write(data)
            f.flush()
            os.fsync(f.fileno())

        os.replace(temporary, target)
        return "updated"
    except BaseException:
        try:
            temporary.unlink()
        except FileNotFoundError:
            pass
        raise


def verify_with_command(target: Path, command: list[str]) -> None:
    result = subprocess.run(
        [*command, str(target)],
        capture_output=True,
        text=True,
        timeout=30,
        check=False,
    )

    if result.returncode != 0:
        raise AutomationError(
            json.dumps(
                {
                    "returncode": result.returncode,
                    "stdout": result.stdout[-2000:],
                    "stderr": result.stderr[-2000:],
                },
                ensure_ascii=False,
            )
        )


def run(args: argparse.Namespace, logger: logging.Logger) -> int:
    run_id = str(uuid.uuid4())
    started_at = utc_now()

    logger.info(
        "任务开始",
        extra={"run_id": run_id, "action": "download-publish"},
    )

    try:
        data = download(args.url, args.timeout)
        input_hash = digest_bytes(data)

        logger.info(
            "下载完成",
            extra={"run_id": run_id, "action": "download"},
        )

        result = atomic_publish(args.output, data)

        logger.info(
            "文件发布完成: %s, sha256=%s",
            result,
            input_hash,
            extra={"run_id": run_id, "action": "publish"},
        )

        if args.verify_command:
            verify_with_command(args.output, args.verify_command)
            logger.info(
                "校验命令完成",
                extra={"run_id": run_id, "action": "verify"},
            )

        audit = {
            "run_id": run_id,
            "action": "download-publish",
            "object": str(args.output),
            "input_hash": input_hash,
            "result": result,
            "status": "succeeded",
            "started_at": started_at,
            "finished_at": utc_now(),
        }
        print(json.dumps(audit, ensure_ascii=False))
        return 0

    except RetryableDownloadError as exc:
        logger.warning(
            "下载失败,结果可重试: %s",
            exc,
            extra={"run_id": run_id, "action": "download"},
        )
        print(f"可重试失败: {exc}", file=sys.stderr)
        return  temporary_exit_code()

    except AutomationError as exc:
        logger.error(
            "任务失败: %s",
            exc,
            extra={"run_id": run_id, "action": "run"},
        )
        print(f"任务失败: {exc}", file=sys.stderr)
        return 1

    except TimeoutError as exc:
        logger.exception(
            "任务结果未知: %s",
            exc,
            extra={"run_id": run_id, "action": "run"},
        )
        print("任务结果未知,请先核验外部状态", file=sys.stderr)
        return 2


def temporary_exit_code() -> int:
    # 约定:75 表示暂时性失败,供调度器决定是否重试。
    return 75


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(prog="wr-auto")
    subparsers = parser.add_subparsers(required=True)

    publish = subparsers.add_parser("publish")
    publish.add_argument("url")
    publish.add_argument("output", type=Path)
    publish.add_argument("--timeout", type=float, default=15)
    publish.add_argument(
        "--verify-command",
        nargs="+",
        help="校验命令,例如 sha256sum",
    )
    publish.set_defaults(func=run)

    return parser


def main() -> int:
    parser = build_parser()
    args = parser.parse_args()
    logger = logging.getLogger("wr_automation")
    logging.basicConfig(
        level=logging.INFO,
        format="%(asctime)s %(levelname)s %(message)s",
    )
    return args.func(args, logger)


if __name__ == "__main__":
    raise SystemExit(main())

运行形式:

python wr_auto.py publish \
  https://example.com/data.json \
  workspace/output/data.json

带校验命令:

python wr_auto.py publish \
  https://example.com/data.json \
  workspace/output/data.json \
  --verify-command python -m json.tool

这段程序的完整数据流是:

flowchart TD
    A[命令行参数] --> B[构造 Request]
    B --> C{下载}
    C -->|网络暂时失败| D[返回 75]
    C -->|HTTP 明确失败| E[返回 1]
    C -->|成功| F[读取完整字节]
    F --> G[计算 SHA-256]
    G --> H[写入同目录临时文件]
    H --> I[flush + fsync]
    I --> J[os.replace 原子发布]
    J --> K{可选校验命令}
    K -->|退出码 0| L[审计 succeeded]
    K -->|非零或超时| M[审计 failed/unknown]

几个关键点不能省略:

  1. 先完整下载,再发布。目标文件不会在下载过程中被覆盖。
  2. 先写临时文件,再替换目标。读取者不会看到下载中的目标文件。
  3. 使用内容哈希。可以判断重复下载是否产生了相同内容。
  4. 校验命令放在发布之后并不总是合适。如果校验失败,目标文件已经替换,恢复策略应明确:保留失败版本、回滚旧版本,还是把目标标记为不可用。
  5. 退出码 75 只是应用约定。脚本和调度器必须共同定义其含义,不能把它当作 Python 或操作系统的普遍规范。
  6. 命令超时不等于发布回滚。如果校验命令本身有外部副作用,仍需按未知结果处理。

九、并发和定时执行:重复启动不是异常,而是常态

定时任务可能出现:

第 1 轮任务执行时间超过调度周期
第 2 轮任务启动
两个实例同时修改同一文件

最简单的单机锁可以利用目录创建的“已存在则失败”行为:

from contextlib import contextmanager
from pathlib import Path
import shutil


@contextmanager
def directory_lock(lock_path: Path):
    try:
        lock_path.mkdir()
    except FileExistsError as exc:
        raise RuntimeError("已有任务运行中") from exc

    try:
        yield
    finally:
        shutil.rmtree(lock_path, ignore_errors=True)

但目录锁有一个关键问题:进程崩溃后锁目录可能残留。加入 PID 和时间戳只能帮助诊断,不能可靠判断持有者是否仍然存活。生产系统通常需要:

  • 操作系统文件锁;
  • 数据库租约;
  • 带过期时间的分布式锁;
  • 调度器本身提供的并发限制;
  • 人工确认后清理陈旧锁。

锁只解决“谁可以同时执行”,不解决“执行到哪一步”。因此锁和幂等状态应配合使用:

取得锁
读取任务状态
若已成功且输入哈希相同:退出
若运行中且租约有效:退出
若运行中但租约过期:接管
执行任务
更新成功、失败或未知
释放锁

如果释放锁的 finally 在进程被强制终止时没有执行,仍然必须依靠租约或人工恢复机制。


十、故障路径:按“已知事实”而不是按异常名称恢复

10.1 文件写入失败

创建临时文件
写入中断

正确恢复:

  • 删除临时文件;
  • 保留旧目标;
  • 审计为 failed
  • 下次可以重新执行。

10.2 原子替换失败

临时文件完整
os.replace 失败

正确恢复:

  • 临时文件仍可能存在;
  • 目标文件通常仍是旧版本;
  • 记录源路径、目标路径和操作系统错误;
  • 不要直接认为目标已经更新。

10.3 命令返回非零

命令启动成功
命令执行
返回码 2

正确恢复取决于命令协议:

  • 参数错误:修复输入,不重试;
  • 临时网络失败:允许重试;
  • 部分完成:先查询外部状态;
  • 不可逆操作:进入人工或补偿流程。

10.4 网络超时

请求发出
客户端超时

此时结果可能是:

A. 请求未到达服务端
B. 服务端收到但未处理
C. 服务端已处理,响应丢失
D. 服务端处理失败,但客户端不知道

因此,对写操作的恢复顺序应是:

  1. 使用幂等键查询;
  2. 查询确定性资源名;
  3. 用内容哈希校验;
  4. 确认未提交后再重试;
  5. 无法确认时记录 unknown,不要盲目重复。

十一、测试:把时间、网络、命令和文件系统替换成可控对象

自动化脚本难测,通常不是因为业务逻辑复杂,而是因为代码直接绑定了:

  • time.sleep
  • urlopen
  • subprocess.run
  • 当前工作目录;
  • 真实网络;
  • 真实系统命令。

将这些依赖作为参数或封装在小函数中,就可以测试状态转换。

例如测试原子发布:

from pathlib import Path


def test_publish_is_idempotent(tmp_path: Path):
    target = tmp_path / "out.txt"

    first = atomic_publish(target, b"hello")
    second = atomic_publish(target, b"hello")

    assert first == "updated"
    assert second == "unchanged"
    assert target.read_bytes() == b"hello"

测试不同内容:

def test_publish_updates_changed_content(tmp_path: Path):
    target = tmp_path / "out.txt"

    atomic_publish(target, b"old")
    result = atomic_publish(target, b"new")

    assert result == "updated"
    assert target.read_bytes() == b"new"

测试命令失败:

import subprocess
import pytest


def test_verify_command_failure(tmp_path: Path):
    target = tmp_path / "data.txt"
    target.write_text("data", encoding="utf-8")

    with pytest.raises(AutomationError):
        verify_with_command(
            target,
            [sys.executable, "-c", "raise SystemExit(3)"],
        )

如果项目只使用标准库,可以使用 unittestunittest.mock,不需要把网络和外部命令带入单元测试。真实网络、真实命令和真实文件系统则属于集成测试,应明确标记并单独运行。


十二、常见错误和边界

把“存在”当成“已成功”

if target.exists():
    return

文件存在可能意味着:

  • 上一次执行已经成功;
  • 上一次执行留下半成品;
  • 另一个任务正在写;
  • 文件内容属于旧输入;
  • 文件是攻击者放置的替代物。

至少要结合内容哈希、版本号、文件大小、校验结果和审计状态。

把“异常”当成“没有副作用”

网络异常只说明调用方没有得到预期响应,不说明远端没有发生变化。

shell=False 当成完整安全措施

参数数组可以避免常见 shell 注入,但仍需考虑:

  • 可执行文件解析路径;
  • 环境变量;
  • 工作目录;
  • 输入文件权限;
  • 外部程序自身的解析漏洞;
  • Windows 批处理文件的特殊行为。

官方文档还指出,Windows 的批处理文件可能由操作系统通过 shell 启动,即使调用方没有显式传入 shell=True。(docs.python.org)

fsync() 当成全系统事务提交

fsync() 处理文件描述符对应的文件刷新,不会让网络服务、数据库事务和审计记录自动一起提交。一个文件已经持久化,并不代表远端状态和审计状态也已经持久化。

把日志当成审计数据库

日志适合观察和诊断,不适合直接承担幂等状态机。需要查询“某个业务键是否已经成功”时,应使用结构化持久化状态,例如 SQLite、数据库表或外部任务系统。

无限重试

无限重试会把暂时性故障变成资源耗尽:

网络不可用
任务无限重试
线程、进程、连接和日志不断增长

重试必须有上限;超过上限后应输出可定位的失败状态,并交给调度器、补偿任务或人工处理。


十三、可靠自动化的最小闭环

一个文件、命令和网络混合的自动化任务,至少应形成如下闭环:

确定输入
  ↓
生成 run_id 和幂等键
  ↓
记录 started
  ↓
执行有超时的读取或命令
  ↓
区分成功、明确失败、未知结果
  ↓
文件使用临时文件和原子替换
  ↓
外部副作用使用幂等键或可查询结果
  ↓
记录 succeeded / failed / unknown
  ↓
通过退出码通知调度器

可以用下面的条件判断设计是否足够稳健:

文件条件

发布成功临时文件完整校验通过替换成功\text{发布成功} \Rightarrow \text{临时文件完整} \land \text{校验通过} \land \text{替换成功}

重试条件

允许重试错误暂时性副作用可重复或可查询未超出时间与次数预算\text{允许重试} \Rightarrow \text{错误暂时性} \land \text{副作用可重复或可查询} \land \text{未超出时间与次数预算}

幂等条件

重复执行安全f(f(x))=f(x)\text{重复执行安全} \Rightarrow f(f(x)) = f(x)

或者,若原始操作本身不是幂等的,则必须引入幂等键、去重记录或确定性资源标识,将重复请求映射到同一个业务结果。

审计条件

可追溯主体时间对象输入尝试次数最终状态\text{可追溯} \Rightarrow \text{主体} \land \text{时间} \land \text{对象} \land \text{输入} \land \text{尝试次数} \land \text{最终状态}

自动化脚本真正可靠的标志,不是“正常时能跑通”,而是遇到超时、崩溃、重复启动、部分写入、命令失败和网络断开后,系统仍能回答三个问题:

  1. 现在已经发生了什么?
  2. 再执行一次会不会造成重复副作用?
  3. 如果结果未知,如何验证并恢复?

当文件提交、命令执行、网络请求、重试、幂等和审计都围绕这三个问题设计时,脚本才从一次性代码变成了可运行、可恢复、可验证的自动化工具。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。