1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
| // 事件存储
interface EventStore {
append(aggregateId: string, events: Event[]): Promise<void>;
getEvents(aggregateId: string): Promise<Event[]>;
getEventsFromVersion(aggregateId string, version: number): Promise<Event[]>;
}
// 实现
class PostgresEventStore implements EventStore {
constructor(private db: Database) {}
async append(aggregateId: string, events: Event[]): Promise<void> {
await this.db.transaction(async (tx) => {
// 检查版本冲突(乐观锁)
const current = await this.getCurrentVersion(tx, aggregateId);
if (current !== events[0].aggregateVersion - 1) {
throw new ConcurrencyError();
}
// 批量插入事件
for (const event of events) {
await tx.insert('events', {
id: event.eventId,
aggregate_id: aggregateId,
event_type: event.eventType,
data: JSON.stringify(event.data),
version: event.aggregateVersion,
timestamp: event.timestamp
});
}
});
}
async getEvents(aggregateId: string): Promise<Event[]> {
const rows = await this.db.query(
`SELECT * FROM events
WHERE aggregate_id = $1
ORDER BY version ASC`,
[aggregateId]
);
return rows.map(row => ({
eventId: row.id,
eventType: row.event_type,
data: JSON.parse(row.data),
aggregateVersion: row.version,
timestamp: row.timestamp
}));
}
// 从事件重建聚合状态
async rebuildAggregate<T>(
aggregateId: string,
initialState: T,
applyEvent: (state: T, event: Event) => T
): Promise<T> {
const events = await this.getEvents(aggregateId);
return events.reduce(
(state, event) => applyEvent(state, event),
initialState
);
}
}
// 使用示例:订单聚合
interface OrderState {
id: string;
status: 'created' | 'paid' | 'shipped' | 'cancelled';
items: OrderItem[];
total: number;
}
function applyOrderEvent(state: OrderState, event: Event): OrderState {
switch (event.eventType) {
case 'order.created':
return {
...state,
id: event.aggregateId,
status: 'created',
items: event.data.items,
total: event.data.total
};
case 'order.paid':
return {
...state,
status: 'paid'
};
case 'order.shipped':
return {
...state,
status: 'shipped'
};
default:
return state;
}
}
// 重建订单状态
const orderState = await eventStore.rebuildAggregate(
orderId,
{} as OrderState,
applyOrderEvent
);
|