PHP 怎么事件存储

wen PHP项目 2

本文目录导读:

PHP 怎么事件存储

  1. 数据库表设计(关系型)
  2. 使用 PDO 实现事件存储类
  3. 领域事件抽象类
  4. 使用 NoSQL(MongoDB)
  5. 使用第三方库(推荐)
  6. 关键设计要点
  7. 使用示例
  8. 注意事项

在 PHP 中实现事件存储(Event Sourcing 的核心组件),通常有以下几种方案,我会从简单到复杂,给出实际可用的代码示例。

数据库表设计(关系型)

最基础的事件存储使用 MySQL/PostgreSQL:

CREATE TABLE events (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    aggregate_id VARCHAR(36) NOT NULL,
    aggregate_type VARCHAR(100) NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    payload JSON NOT NULL,
    metadata JSON,
    version INT NOT NULL,
    created_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
    INDEX idx_aggregate (aggregate_id, aggregate_type, version),
    UNIQUE KEY uk_aggregate_version (aggregate_id, aggregate_type, version)
) ENGINE=InnoDB;

使用 PDO 实现事件存储类

<?php
class EventStore {
    private PDO $pdo;
    public function __construct(PDO $pdo) {
        $this->pdo = $pdo;
    }
    // 追加事件
    public function append(DomainEvent $event, int $expectedVersion = null): void {
        $this->pdo->beginTransaction();
        try {
            // 乐观锁检查(可选)
            if ($expectedVersion !== null) {
                $stmt = $this->pdo->prepare(
                    "SELECT COUNT(*) FROM events 
                     WHERE aggregate_id = :id AND version = :version"
                );
                $stmt->execute([
                    ':id' => $event->getAggregateId(),
                    ':version' => $expectedVersion
                ]);
                if ($stmt->fetchColumn() == 0) {
                    throw new ConcurrencyException("事件版本冲突");
                }
            }
            // 插入事件
            $stmt = $this->pdo->prepare(
                "INSERT INTO events 
                 (aggregate_id, aggregate_type, event_type, payload, metadata, version)
                 VALUES (:aggregate_id, :aggregate_type, :event_type, :payload, :metadata, :version)"
            );
            $stmt->execute([
                ':aggregate_id' => $event->getAggregateId(),
                ':aggregate_type' => $event->getAggregateType(),
                ':event_type' => get_class($event),
                ':payload' => json_encode($event->toArray(), JSON_THROW_ON_ERROR),
                ':metadata' => json_encode($event->getMetadata() ?? []),
                ':version' => $event->getVersion()
            ]);
            $this->pdo->commit();
        } catch (PDOException $e) {
            $this->pdo->rollBack();
            // 检查唯一约束违规(并发冲突)
            if ($e->getCode() == '23000') {
                throw new ConcurrencyException("并发写入冲突");
            }
            throw $e;
        }
    }
    // 按聚合ID读取事件流
    public function loadStream(string $aggregateId, string $aggregateType, int $fromVersion = 1): array {
        $stmt = $this->pdo->prepare(
            "SELECT * FROM events 
             WHERE aggregate_id = :id AND aggregate_type = :type AND version >= :version
             ORDER BY version ASC"
        );
        $stmt->execute([
            ':id' => $aggregateId,
            ':type' => $aggregateType,
            ':version' => $fromVersion
        ]);
        return $this->hydrateEvents($stmt->fetchAll());
    }
    // 按事件类型读取(用于投影)
    public function loadByEventType(string $eventType, array $metadataFilter = []): array {
        $sql = "SELECT * FROM events WHERE event_type = :type";
        $params = [':type' => $eventType];
        if (!empty($metadataFilter)) {
            $sql .= " AND metadata->> :key = :value";
        }
        $stmt = $this->pdo->prepare($sql);
        $stmt->execute($params);
        return $this->hydrateEvents($stmt->fetchAll());
    }
    // 事件流迭代(支持大数据量)
    public function iterateStream(string $aggregateId): \Generator {
        $stmt = $this->pdo->prepare(
            "SELECT * FROM events 
             WHERE aggregate_id = :id ORDER BY version ASC"
        );
        $stmt->execute([':id' => $aggregateId]);
        while ($row = $stmt->fetch(PDO::FETCH_ASSOC)) {
            yield $this->hydrateEvent($row);
        }
    }
    // 将数据库行转换为领域事件对象
    private function hydrateEvents(array $rows): array {
        return array_map([$this, 'hydrateEvent'], $rows);
    }
    private function hydrateEvent(array $row): DomainEvent {
        $payload = json_decode($row['payload'], true);
        // 假设事件类有静态的 fromArray 方法重构
        $eventClass = $row['event_type'];
        return $eventClass::fromArray($payload);
    }
    // 获取聚合的最新版本(用于快照)
    public function getLatestVersion(string $aggregateId): int {
        $stmt = $this->pdo->prepare(
            "SELECT MAX(version) FROM events WHERE aggregate_id = :id"
        );
        $stmt->execute([':id' => $aggregateId]);
        return (int) $stmt->fetchColumn() ?? 0;
    }
}

领域事件抽象类

<?php
abstract class DomainEvent {
    private string $aggregateId;
    private string $aggregateType;
    private int $version;
    private array $metadata = [];
    private \DateTimeImmutable $occurredAt;
    public function __construct(
        string $aggregateId, 
        string $aggregateType, 
        int $version,
        array $metadata = []
    ) {
        $this->aggregateId = $aggregateId;
        $this->aggregateType = $aggregateType;
        $this->version = $version;
        $this->metadata = $metadata;
        $this->occurredAt = new \DateTimeImmutable();
    }
    // 将事件转换为数组(持久化用)
    abstract public function toArray(): array;
    // 从数组重构事件(读取用)
    abstract public static function fromArray(array $data): self;
    // Getters...
}
// 具体事件示例
class OrderCreated extends DomainEvent {
    private string $customerId;
    private float $total;
    public function __construct(string $aggregateId, string $customerId, float $total, int $version) {
        parent::__construct($aggregateId, 'order', $version);
        $this->customerId = $customerId;
        $this->total = $total;
    }
    public function toArray(): array {
        return [
            'customerId' => $this->customerId,
            'total' => $this->total
        ];
    }
    public static function fromArray(array $data): self {
        return new self(
            $data['aggregateId'] ?? '',
            $data['customerId'],
            $data['total'],
            $data['version'] ?? 1
        );
    }
    public function getCustomerId(): string {
        return $this->customerId;
    }
    public function getTotal(): float {
        return $this->total;
    }
}

使用 NoSQL(MongoDB)

<?php
class MongoEventStore {
    private MongoDB\Collection $collection;
    public function __construct(MongoDB\Collection $collection) {
        $this->collection = $collection;
    }
    public function append(DomainEvent $event): void {
        $this->collection->insertOne([
            'aggregate_id' => $event->getAggregateId(),
            'aggregate_type' => $event->getAggregateType(),
            'event_type' => get_class($event),
            'payload' => $event->toArray(),
            'metadata' => $event->getMetadata(),
            'version' => $event->getVersion(),
            'created_at' => new MongoDB\BSON\UTCDateTime()
        ]);
    }
    public function loadStream(string $aggregateId): array {
        $cursor = $this->collection->find(
            ['aggregate_id' => $aggregateId],
            ['sort' => ['version' => 1]]
        );
        $events = [];
        foreach ($cursor as $doc) {
            $eventClass = $doc['event_type'];
            $events[] = $eventClass::fromArray($doc['payload']);
        }
        return $events;
    }
}

使用第三方库(推荐)

对于生产环境,建议使用成熟的库:

EventSauce

composer require eventsauce/event-sourcing
use EventSauce\EventSourcing\DefaultEventSourcingRepository;
use EventSauce\EventSourcing\MessageRepository;
use EventSauce\EventSourcing\Serialization\ConstructingMessageSerializer;
// 定义聚合
class Order extends AggregateRoot {
    private function applyOrderCreated(OrderCreated $event): void {
        $this->customerId = $event->customerId();
    }
    public function create(string $customerId): static {
        $this->recordThat(new OrderCreated($customerId));
        return $this;
    }
}
// 使用
$repository = new DefaultEventSourcingRepository(
    $aggregateRootRepository,
    new ConstructingMessageSerializer()
);
$order = Order::create('customer-123');
$repository->persist($order->aggregateRootId(), $order->releaseEvents());

prooph/event-sourcing

composer require prooph/event-sourcing prooph/event-store-pdo

关键设计要点

// 1. 乐观并发控制
interface OptimisticLocking {
    public function getVersion(): int;
}
// 2. 事件序列化
trait EventSerializable {
    public function serialize(): string {
        return json_encode($this->toArray());
    }
    public static function deserialize(string $data): self {
        return static::fromArray(json_decode($data, true));
    }
}
// 3. 快照支持(可选优化)
class SnapshotStore {
    // 定期保存聚合状态,减少事件回溯
    public function save(string $aggregateId, int $version, array $state): void {
        // 保存快照
    }
    public function load(string $aggregateId): ?Snapshot {
        // 加载最近的快照
    }
}

使用示例

// 初始化
$pdo = new PDO('mysql:host=localhost;dbname=events', 'user', 'pass');
$eventStore = new EventStore($pdo);
// 创建聚合根
$order = new Order();
$order->create('customer-01', 250.00);
$order->complete();
// 持久化事件
$events = $order->releaseEvents();
foreach ($events as $event) {
    $eventStore->append($event);
}
// 重新构建聚合
$stream = $eventStore->loadStream('order-123', 'order');
$order = Order::reconstituteFromHistory($stream);

注意事项

  1. 事务性:事件存储通常与投影(Projection)不同步,需要处理最终一致性
  2. 事件版本化:事件结构可能变化,需要在事件中注入版本号
  3. 测试:确保有测试环境的独立事件流
  4. 监控:记录事件写入延迟和失败率
  5. 归档:旧事件可以归档到大数据存储(如 S3)

选择哪种方案取决于你的需求:如果追求简单有效,使用 PDO + MySQL 是很好的选择;如果项目复杂,考虑使用 EventSauce 这类成熟库。

抱歉,评论功能暂时关闭!