0Pricing
PHP Academy · 课时

构建事件驱动的工作流

通过事件和幂等性协调服务

构建事件驱动的工作流 是 CoddyKit 上的免费 PHP Academy 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 PHP Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 PHP Academy 课程共包含 4 节课。

事件驱动的工作流

单条消息很简单。一个工作流——“下单 → 预留库存 → 扣款 → 发货 → 通知”,跨越多个服务——正是事件驱动设计发挥价值的地方;如果处理不当,也正是问题开始出现的地方。

本课将介绍事件编舞与编排、用于原子发布的发件箱模式、用于分布式回滚的长事务,以及将这一切连接起来的幂等性。

编舞与编排

协调多步骤流程的两种方式:

  • 事件编舞 — 每个服务响应事件并发出自己的事件;没有中央大脑。耦合松散,但整体流程是隐式的,难以追踪。
  • 编排 — 中央协调器告诉每个服务下一步该做什么。流程明确且可观测,但编排器会成为耦合点。

经验法则:简单扇出使用事件编舞;当流程包含许多有序步骤且需要可见状态时,使用编排。

事件与命令

请有意地命名消息:

  • 事件陈述过去发生的事实:OrderPlaced。发布者不关心谁在监听。
  • 命令请求特定处理程序在未来执行某项操作:ChargeCard。

事件驱动事件编舞;命令驱动编排。混用词汇(表面上是“事件”,却暗中要求由一个处理程序处理)是产生隐性耦合的常见原因。

<?php
final class OrderPlaced {
    public function __construct(
        public readonly string $orderId,
        public readonly string $customerId,
        public readonly int $amountCents,
        public readonly string $occurredAt,
    ) {}
}

$e = new OrderPlaced('o-42', 'c-7', 1990, gmdate('c'));
echo json_encode($e), "\n";

双写问题

典型错误是:处理程序将数据库更新和消息发布作为两个独立操作。如果进程在两者之间退出,就会产生不一致——数据行已更改却没有发送事件,或者反过来。

<?php
// BROKEN: not atomic. A crash between the two lines corrupts state.
function placeOrder(PDO $db, $broker, array $o): void {
    $db->prepare('INSERT INTO orders ...')->execute($o);
    // <-- crash here = row exists but no event ever published
    $broker->publish('OrderPlaced', json_encode($o));
}

事务性发件箱

解决办法是发件箱模式:在更改数据的同一个数据库事务中,将事件插入 outbox 表。独立的中继进程读取尚未发布的行,并将其发送到消息代理。一次原子提交,避免双写。

<?php
function placeOrder(PDO $db, array $o): void {
    $db->beginTransaction();
    $db->prepare('INSERT INTO orders (id, total) VALUES (?, ?)')
       ->execute([$o['id'], $o['total']]);
    // Same transaction -> atomic with the business write
    $db->prepare('INSERT INTO outbox (id, type, payload) VALUES (?, ?, ?)')
       ->execute([bin2hex(random_bytes(8)), 'OrderPlaced', json_encode($o)]);
    $db->commit();
}

中继(发布者)

工作进程会轮询发件箱(或通过 CDC 跟踪数据库变更日志),发布每一行,然后将其标记为已发送。由于中继可能在发布之后、标记之前崩溃,因此它本身也是至少一次投递——这没有问题,因为消费者具有幂等性。

<?php
function relayOutbox(PDO $db, $broker): void {
    $rows = $db->query(
        'SELECT id, type, payload FROM outbox
         WHERE published_at IS NULL ORDER BY created_at LIMIT 100'
    )->fetchAll(PDO::FETCH_ASSOC);

    foreach ($rows as $r) {
        $broker->publish($r['type'], $r['payload'], messageId: $r['id']);
        $db->prepare('UPDATE outbox SET published_at = now() WHERE id = ?')
           ->execute([$r['id']]);
    }
}

幂等消费者(再谈)

由于中继和消息代理都采用至少一次投递,下游处理程序一定会看到重复消息。每个消费者都会记录已处理的消息标识,并对重复消息直接跳过——这就是您之前学过的同一个去重门禁,现在应用于每个服务。

<?php
function onOrderPlaced(PDO $db, string $messageId, array $data): void {
    $db->beginTransaction();
    try {
        $db->prepare('INSERT INTO inbox (message_id) VALUES (?)')
           ->execute([$messageId]); // unique index = dedup
    } catch (PDOException $e) {
        $db->rollBack();
        return; // already handled this message
    }
    reserveStock($data['orderId']);
    $db->commit();
}
function reserveStock(string $id): void {}

长事务:分布式回滚

您无法跨多个服务开启一个 ACID 事务。长事务会将长时间运行的流程建模为一系列本地事务,每个事务都有一个用于撤销自身的补偿操作。如果第 3 步失败,就按相反顺序执行第 2 步和第 1 步的补偿操作。

例如:库存已经预留,但支付失败 → 发出 ReleaseStock 进行补偿。不存在自动回滚——您必须为每一步设计撤销操作。

编排式长事务

编排器驱动长事务:成功时推进流程,失败时分派补偿操作。请持久化长事务的状态,以便发生崩溃后能够恢复。

<?php
function handleStepResult(array $saga, string $step, bool $ok, $bus): array {
    if ($ok) {
        $next = ['reserveStock' => 'chargeCard', 'chargeCard' => 'ship'][$step] ?? null;
        if ($next) { $bus->send($next, $saga['orderId']); $saga['state'] = $next; }
        else { $saga['state'] = 'completed'; }
    } else {
        // Run compensations in reverse for whatever already succeeded
        foreach (array_reverse($saga['done']) as $s) {
            $bus->send('compensate.' . $s, $saga['orderId']);
        }
        $saga['state'] = 'compensating';
    }
    return $saga;
}

长流程中的超时

长事务中的某一步可能永远不返回——支付服务停机,或者人工审批始终没有完成。如果没有超时,长事务会一直挂起并占用预留资源。请为每一步持久化截止时间;调度器扫描超期的长事务,并触发失败/补偿路径。

<?php
function reapTimedOutSagas(PDO $db, $bus): void {
    $rows = $db->query(
        "SELECT order_id, state FROM sagas
         WHERE state NOT IN ('completed','compensating')
           AND deadline_at < now()"
    )->fetchAll(PDO::FETCH_ASSOC);

    foreach ($rows as $r) {
        echo "Saga {$r['order_id']} timed out at step {$r['state']}\n";
        $bus->send('saga.compensate', $r['order_id']); // trigger rollback
    }
}

版本管理与可观测性

工作流会持续运行数年;事件必须安全地演进:

  • 为每个事件添加 version(或模式);消费者应兼容未知的新字段,绝不能假定字段一定存在。
  • 优先采用增量式变更;绝不要改变现有字段的原有含义。
  • 在每条消息中传递一个关联 ID,这样您就能在所有服务的日志和追踪记录中追溯同一笔业务事务。

没有关联 ID,几乎不可能调试跨越五个服务的编舞式流程。

<?php
$envelope = [
    'type'          => 'OrderPlaced',
    'version'       => 2,
    'correlationId' => $incoming['correlationId'] ?? bin2hex(random_bytes(8)),
    'occurredAt'    => gmdate('c'),
    'data'          => ['orderId' => 'o-42'],
];
echo json_encode($envelope, JSON_PRETTY_PRINT), "\n";

快速检查

避免双重写入问题。

回顾

现在您已经能够设计可靠的事件驱动工作流:

  • 根据流程复杂度,为每个流程选择事件编舞(事件)或命令编排(命令)。
  • 使用事务发件箱加中继解决双重写入问题。
  • 通过收件箱/去重键,让每个消费者都具备幂等性。
  • 使用分布式长事务模式,通过补偿操作实现分布式回滚。
  • 以增量方式演进事件,并贯穿传递关联 ID以便追踪。

这些模式能将松散的消息转变为可靠且可观测的业务流程。

常见问题解答

「构建事件驱动的工作流」课时是免费的吗?

是的 — 「构建事件驱动的工作流」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 PHP Academy 课程的其余内容,请升级到 CoddyKit PRO。 PHP Academy 课程共包含 4 节课。

「构建事件驱动的工作流」这节课中我会学到什么?

通过事件和幂等性协调服务 你通过在浏览器中直接运行的动手代码来练习 PHP Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 PHP Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 PHP Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。

「构建事件驱动的工作流」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 PHP Academy 课中编写并运行代码吗?

能。每节 PHP Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 异步消息为何重要
  2. 在 PHP 中使用 RabbitMQ
  3. 使用 PHP 的 Apache Kafka
  4. 构建事件驱动的工作流
← 返回 PHP Academy