本文目录导读:

在 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);
注意事项
- 事务性:事件存储通常与投影(Projection)不同步,需要处理最终一致性
- 事件版本化:事件结构可能变化,需要在事件中注入版本号
- 测试:确保有测试环境的独立事件流
- 监控:记录事件写入延迟和失败率
- 归档:旧事件可以归档到大数据存储(如 S3)
选择哪种方案取决于你的需求:如果追求简单有效,使用 PDO + MySQL 是很好的选择;如果项目复杂,考虑使用 EventSauce 这类成熟库。