用 Codex 或 Claude 帮你安装 复制这段 Prompt,粘贴到 Codex、Claude 或其他助手里,让它检查 Skill 页面并帮你完成安装。
直接命令不会经过审查 Prompt;运行前请先检查来源。
npx skills add https://github.com/ffsshhttiikk/opencode-agents-skills --skill event-driven命令会保持在同一行。复制前请横向滚动并检查完整内容。
想先保存到本地?可下载 SkillsMP 当前能够提供的文件。
正在显示 SKILL.md
基于 SOC 职业分类
| name | event-driven |
| description | Event-driven architecture patterns and best practices |
| license | MIT |
| compatibility | opencode |
| metadata | {"audience":"developers","category":"architecture"} |
When designing event-driven architectures or implementing event systems.
┌─────────────────────────────────────────────────────────────────────┐
│ Event Producers │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ Service │ │ Service │ │ Service │ │ Service │ │
│ │ A │───►│ B │───►│ C │───►│ D │ │
│ └─────────┘ └─────────┘ └─────────┘ └─────────┘ │
│ │ │ │ │ │
│ └──────────────┴──────────────┴──────────────┘ │
│ │ │
│ ┌────────▼────────┐ │
│ │ Event Bus │ │
│ │ (Kafka/Rabbit) │ │
│ └────────┬────────┘ │
│ │ │
│ ┌───────────────────┼───────────────────┐ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │Consumer │ │Consumer │ │Consumer │ │
│ │ 1 │ │ 2 │ │ 3 │ │
│ └─────────┘ └─────────┘ └─────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────┘
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any, Dict
import json
import uuid
@dataclass
class Event:
"""Base event structure."""
event_id: str = field(default_factory=lambda: str(uuid.uuid4()))
event_type: str = ""
aggregate_id: str = ""
aggregate_type: str = ""
occurred_at: datetime = field(default_factory=datetime.utcnow)
payload: Dict[str, Any] = field(default_factory=dict)
metadata: Dict[str, Any] = field(default_factory=dict)
def to_dict(self) -> dict:
return {
"event_id": self.event_id,
"event_type": self.event_type,
"aggregate_id": self.aggregate_id,
"aggregate_type": self.aggregate_type,
"occurred_at": self.occurred_at.isoformat(),
"payload": self.payload,
: .metadata,
}
() -> :
data[] = datetime.fromisoformat(data[])
cls(**data)
() -> :
json.dumps(.to_dict())
() -> :
cls.from_dict(json.loads(json_str))
():
event_type: =
payload: [, ] = field(default_factory=)
():
event_type: =
payload: [, ] = field(default_factory=)
:
TYPE_MAPPING = {
: UserCreatedEvent,
: OrderPlacedEvent,
}
() -> Event:
event_type = data.get()
event_type cls.TYPE_MAPPING:
event_cls = cls.TYPE_MAPPING[event_type]
event_cls.from_dict(data)
Event.from_dict(data)
from abc import ABC, abstractmethod
from typing import List, Optional
from datetime import datetime
class Aggregate(ABC):
"""Base aggregate with event sourcing."""
def __init__(self, aggregate_id: str) -> None:
self.id = aggregate_id
self._events: List[Event] = []
self._version: int = 0
def apply(self, event: Event) -> None:
"""Apply event to aggregate."""
self._events.append(event)
self._version += 1
self._handle_event(event)
@abstractmethod
def _handle_event(self, event: Event) -> None:
"""Handle event (to be implemented by subclass)."""
pass
def get_pending_events(self) -> List[Event]:
"""Get events not yet persisted."""
return self._events
() -> :
._events = []
() -> :
aggregate = cls(aggregate_id)
event events:
aggregate.apply(event)
aggregate
():
() -> :
().__init__(user_id)
.email: [] =
.name: [] =
.created_at: [datetime] =
.updated_at: [datetime] =
() -> :
event = UserCreatedEvent(
aggregate_id=.,
aggregate_type=,
payload={
: email,
: name,
},
)
.apply(event)
() -> :
(event, UserCreatedEvent):
.email = event.payload[]
.name = event.payload[]
.created_at = event.occurred_at
from dataclasses import dataclass
from typing import Dict, Any, Callable
from enum import Enum
import asyncio
class SagaStepStatus(Enum):
PENDING = "pending"
IN_PROGRESS = "in_progress"
COMPLETED = "completed"
FAILED = "failed"
COMPENSATING = "compensating"
@dataclass
class SagaContext:
"""Context passed through saga steps."""
data: Dict[str, Any] = None
status: SagaStepStatus = SagaStepStatus.PENDING
compensating: bool = False
def __post_init__(self):
if self.data is None:
self.data = {}
class Saga:
"""Base saga with compensation support."""
def __init__(self, saga_id: str) -> None:
self.saga_id = saga_id
self.steps: List[Callable] = []
self.compensation_steps: [] = []
.context = SagaContext()
() -> :
.steps.append(forward)
.compensation_steps.append(backward)
() -> :
executed_steps = []
:
step .steps:
success = step(.context)
success:
SagaStepFailed()
executed_steps.append(step)
.context.status = SagaStepStatus.COMPLETED
SagaStepFailed:
.context.status = SagaStepStatus.COMPENSATING
.context.compensating =
step (executed_steps):
index = .steps.index(step)
compensation = .compensation_steps[index]
compensation(.context)
():
() -> :
().__init__(order_id)
.add_step(
.reserve_inventory,
.release_inventory
)
.add_step(
.process_payment,
.refund_payment
)
.add_step(
.create_shipment,
.cancel_shipment
)
() -> :
ctx.data[] =
ctx.data[] =
() -> :
ctx.data.get():
ctx.data[] =
() -> :
ctx.data[] =
ctx.data[] =
() -> :
ctx.data.get():
ctx.data[] =
from abc import ABC, abstractmethod
from typing import List, Optional
import asyncio
class EventConsumer(ABC):
"""Base event consumer."""
def __init__(self, consumer_group: str) -> None:
self.consumer_group = consumer_group
self.running = False
@abstractmethod
async def handle_event(self, event: Event) -> bool:
"""Handle single event. Return True if processed successfully."""
pass
async def process_batch(self, events: List[Event]) -> int:
"""Process batch of events."""
processed = 0
for event in events:
try:
success = await self.handle_event(event)
if success:
processed += 1
except Exception as e:
print(f"Error processing event: {e}")
processed
:
() -> :
.max_retries = max_retries
.dead_letter_queue = []
() -> :
event.metadata[] = event.metadata.get(, ) +
event.metadata[] < .max_retries:
.schedule_retry(event)
:
.dead_letter_queue.append({
: event,
: (error),
: datetime.utcnow().isoformat(),
})
() -> :
delay = ( ** event.metadata[], )
# Ordering guarantees
# Per-partition ordering (Kafka)
# - Events with same key go to same partition
# - Consumer processes in order
# Global ordering
# - Single partition (lower throughput)
# - Or accept eventual ordering
# Key-based partitioning
def get_partition_key(event: Event) -> str:
"""Get partition key for event."""
return f"{event.aggregate_type}:{event.aggregate_id}"
1. Design events as facts
- Immutable, never modify
- Past tense naming: UserCreated, not CreateUser
2. Include enough context
- Sufficient data for consumers
- Avoid frequent lookups
3. Version your events
- Backward compatibility
- Schema evolution
4. Handle idempotency
- Deduplicate events
- Process exactly once
5. Monitor lag
- Track consumer lag
- Alert on delays
6. Test thoroughly
- Event handlers
- Saga rollbacks
7. Document event flow
- Event catalog
- Consumer dependencies