Token导航 LogoToken导航TokenDH.com
开发需要联网github未标认证来源可访问clear审计异常

outbox-pattern发件箱模式

Agent Skill

outbox-pattern 用于处理 GitHub 仓库、Issue、Pull Request 和代码协作信息,适合在 Codex、Claude、Cursor、Gemini CLI 中需要围绕仓库状态、代码变更或协作事项进行整理时使用。可结合来源仓库、安装命令和原始 README 继续核验具体用法。安装前建议确认权限范围、维护状态,以及是否会触发联网、命令执行或文件读写。

总安装

282

周安装

12

GitHub Stars

161

下载量

99
CodexClaudeCursorGemini CLI

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

GitHub

来源数

3

许可证

MIT

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

复制提示词发给支持本地命令或 Skills 的 AI 助手,先确认命令和权限,再让它执行。

请帮我安装这个 Agent Skill:outbox-pattern(发件箱模式)
来源仓库:https://github.com/yonatangross/orchestkit
仓库路径:skills/outbox-pattern
安装命令:
npx skills add https://github.com/yonatangross/orchestkit --skill outbox-pattern
安装前请先检查当前环境是否支持对应 CLI,并向我确认将要执行的命令、安装目录、联网范围和文件读写权限;确认后再执行。

命令行安装

复制命令到本机终端执行。不同来源提供的安装方式可能略有差异;本站展示可直接复制的安装命令,安装前请核对来源页面。

skills.shnpx skills
npx skills add https://github.com/yonatangross/orchestkit --skill outbox-pattern

简介

outbox-pattern 用于处理 GitHub 仓库、Issue、Pull Request 和代码协作信息。

  • 适合在 Codex、Claude、Cursor、Gemini CLI 中围绕仓库状态、代码变更或协作事项进行整理。
  • 通过 npx skills add 命令从 GitHub 仓库安装使用。
  • 安装前需确认权限范围、维护状态,以及是否触发联网、命令执行或文件读写操作。
  • 建议结合原始 README 核验具体用法和功能边界。

SKILL.md

Outbox Pattern ()

Ensure atomic state changes and event publishing by writing both to a database transaction, then publishing asynchronously.

Overview

  • Ensuring database writes and event publishing are atomic
  • Building reliable event-driven microservices
  • Implementing exactly-once message delivery semantics
  • Avoiding dual-write problems (DB + message broker)
  • Decoupling domain logic from message infrastructure
  • High-throughput systems needing CDC-based publishing

Quick Reference

Outbox Table Schema

CREATE TABLE outbox (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    aggregate_type VARCHAR(100) NOT NULL,
    aggregate_id UUID NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    payload JSONB NOT NULL,
    idempotency_key VARCHAR(255) UNIQUE,  -- For consumer deduplication
    created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    published_at TIMESTAMPTZ,
    retry_count INT DEFAULT 0,
    last_error TEXT
);

-- Index for polling unpublished messages
CREATE INDEX idx_outbox_unpublished ON outbox(created_at)
    WHERE published_at IS NULL;

-- Index for aggregate ordering
CREATE INDEX idx_outbox_aggregate ON outbox(aggregate_id, created_at);

-- Index for idempotency key lookups
CREATE INDEX idx_outbox_idempotency ON outbox(idempotency_key)
    WHERE idempotency_key IS NOT NULL;

SQLAlchemy Model

from sqlalchemy.dialects.postgresql import UUID, JSONB
import hashlib

class OutboxMessage(Base):
    __tablename__ = "outbox"

    id = Column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
    aggregate_type = Column(String(100), nullable=False)
    aggregate_id = Column(UUID(as_uuid=True), nullable=False, index=True)
    event_type = Column(String(100), nullable=False)
    payload = Column(JSONB, nullable=False)
    idempotency_key = Column(String(255), unique=True, nullable=True)
    created_at = Column(DateTime, default=lambda: datetime.now(timezone.utc))
    published_at = Column(DateTime, nullable=True)
    retry_count = Column(Integer, default=0)
    last_error = Column(Text, nullable=True)

    @staticmethod
    def generate_idempotency_key(aggregate_id: str, event_type: str, payload: dict) -> str:
        """Generate deterministic idempotency key for deduplication."""
        content = f"{aggregate_id}:{event_type}:{json.dumps(payload, sort_keys=True)}"
        return hashlib.sha256(content.encode()).hexdigest()[:32]

Write to Outbox in Transaction

from sqlalchemy.ext.asyncio import AsyncSession

class OrderService:
    def __init__(self, db: AsyncSession):
        self.db = db

    async def create_order(self, order_data: OrderCreate) -> Order:
        """Create order AND outbox message in single transaction."""
        # Create the order
        order = Order(**order_data.model_dump())
        self.db.add(order)

        # Create outbox message in SAME transaction
        outbox_msg = OutboxMessage(
            aggregate_type="Order",
            aggregate_id=order.id,
            event_type="OrderCreated",
            payload={
                "order_id": str(order.id),
                "customer_id": str(order.customer_id),
                "total": order.total,
            },
            idempotency_key=OutboxMessage.generate_idempotency_key(
                str(order.id), "OrderCreated", {"total": order.total}
            ),
        )
        self.db.add(outbox_msg)

        await self.db.flush()  # Both written atomically
        return order

Polling Publisher

class OutboxPublisher:
    """Polls outbox and publishes to message broker."""

    def __init__(self, session_factory, producer):
        self.session_factory = session_factory
        self.producer = producer

    async def publish_pending(self, batch_size: int = 100) -> int:
        async with self.session_factory() as session:
            stmt = (
                select(OutboxMessage)
                .where(OutboxMessage.published_at.is_(None))
                .order_by(OutboxMessage.created_at)
                .limit(batch_size)
                .with_for_update(skip_locked=True)  # Prevent duplicate processing
            )
            result = await session.execute(stmt)
            messages = result.scalars().all()

            published = 0
            for msg in messages:
                try:
                    await self.producer.publish(
                        topic=f"{msg.aggregate_type.lower()}-events",
                        key=str(msg.aggregate_id),
                        value={
                            "type": msg.event_type,
                            "idempotency_key": msg.idempotency_key,
                            **msg.payload
                        },
                    )
                    msg.published_at = datetime.now(timezone.utc)
                    published += 1
                except Exception as e:
                    msg.retry_count += 1
                    msg.last_error = str(e)

            await session.commit()
            return published

CDC with Debezium (High-Throughput)

# docker-compose.yml - Debezium connector
version: '3.8'
services:
  debezium:
    image: debezium/connect:2.5
    environment:
      BOOTSTRAP_SERVERS: kafka:9092
      GROUP_ID: outbox-connector
      CONFIG_STORAGE_TOPIC: connect-configs
      OFFSET_STORAGE_TOPIC: connect-offsets

# Register outbox connector
# POST http://debezium:8083/connectors
{
  "name": "outbox-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "secret",
    "database.dbname": "app",
    "table.include.list": "public.outbox",
    "transforms": "outbox",
    "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
    "transforms.outbox.table.field.event.id": "id",
    "transforms.outbox.table.field.event.key": "aggregate_id",
    "transforms.outbox.table.field.event.payload": "payload",
    "transforms.outbox.route.topic.replacement": "${routedByValue}-events"
  }
}

Idempotent Consumer

class IdempotentConsumer:
    """Consumer with deduplication using idempotency keys."""

    def __init__(self, db: AsyncSession, redis: Redis):
        self.db = db
        self.redis = redis

    async def process(self, event: dict, handler) -> bool:
        """Process event idempotently - returns False if duplicate."""
        idempotency_key = event.get("idempotency_key")
        if not idempotency_key:
            # No key = always process (but risky)
            await handler(event)
            return True

        # Check Redis cache first (fast path)
        if await self.redis.exists(f"processed:{idempotency_key}"):
            return False  # Already processed

        # Check database (slow path, but durable)
        exists = await self.db.execute(
            select(ProcessedEvent)
            .where(ProcessedEvent.idempotency_key == idempotency_key)
        )
        if exists.scalar_one_or_none():
            # Cache for future fast lookups
            await self.redis.setex(f"processed:{idempotency_key}", 86400, "1")
            return False

        # Process and record
        async with self.db.begin():
            await handler(event)
            self.db.add(ProcessedEvent(idempotency_key=idempotency_key))
            await self.db.flush()

        # Cache the processed key
        await self.redis.setex(f"processed:{idempotency_key}", 86400, "1")
        return True

class ProcessedEvent(Base):
    """Track processed events for idempotency."""
    __tablename__ = "processed_events"

    idempotency_key = Column(String(255), primary_key=True)
    processed_at = Column(DateTime, default=lambda: datetime.now(timezone.utc))

Dapr Outbox Integration

# Using Dapr's built-in outbox support
from dapr.clients import DaprClient

async def create_order_with_dapr(order_data: dict):
    """Dapr handles outbox automatically."""
    async with DaprClient() as client:
        # Single call - Dapr ensures atomicity
        await client.save_state(
            store_name="statestore",
            key=f"order-{order_data['id']}",
            value=order_data,
            state_metadata={
                "outbox.publish": "true",
                "outbox.topic": "orders",
            }
        )

Key Decisions

DecisionOption AOption BRecommendation
DeliveryPollingCDC (Debezium)Polling for simplicity, CDC for > 10K msg/s
Batch sizeSmall (10-50)Large (100-500)Start with 100, tune based on latency
RetryFixed delayExponential backoffExponential with max 5 retries
CleanupDelete publishedArchive to cold storageArchive for audit, delete after 30 days
OrderingPer-aggregateGlobalPer-aggregate via partition key
IdempotencyConsumer-sideBuilt-in keyAlways include idempotency key
ToolCustomDaprDapr if K8s, custom otherwise

Outbox vs CDC Trade-offs

┌────────────────────────────────────────────────────────────────────────┐
│                     POLLING vs CDC COMPARISON                           │
├────────────────────────────────────────────────────────────────────────┤
│                                                                         │
│  POLLING (OutboxPublisher)          CDC (Debezium)                     │
│  ──────────────────────────         ────────────────                   │
│  ✓ Simple to implement              ✓ Higher throughput (100K+ msg/s)  │
│  ✓ No extra infrastructure          ✓ Lower latency (sub-second)       │
│  ✓ Easy to debug                    ✓ No polling overhead              │
│  ✗ Polling overhead                 ✗ Complex infrastructure           │
│  ✗ Higher latency (1-5s)            ✗ Harder to debug                  │
│  ✗ Limited throughput (~10K/s)      ✗ Requires Kafka Connect           │
│                                                                         │
│  USE WHEN:                          USE WHEN:                          │
│  - Starting out                     - High throughput required         │
│  - Simple architecture              - Sub-second latency needed        │
│  - < 10K events/second              - Already using Kafka              │
│                                                                         │
└────────────────────────────────────────────────────────────────────────┘

Anti-Patterns (FORBIDDEN)

# NEVER publish before commit - dual-write problem
await producer.publish(event)  # May succeed but commit may fail!
await session.commit()

# NEVER delete without publishing - events lost
await session.execute(delete(OutboxMessage).where(...))

# NEVER process without locking - causes duplicates
messages = await session.execute(select(OutboxMessage))  # No lock!

# NEVER store large payloads - use URL references instead
OutboxMessage(payload={"file": large_binary_data})

# NEVER ignore ordering - use aggregate_id as partition key
await producer.publish(event_a)  # May arrive out of order!

# NEVER skip idempotency keys
OutboxMessage(payload=event_data)  # No idempotency_key = duplicate risk

# NEVER process without idempotency check on consumer
async def handle(event):
    await process(event)  # Duplicate processing on retry!

Related Skills

  • message-queues - Kafka, RabbitMQ, Redis Streams integration
  • event-sourcing - Full event-sourced architecture patterns
  • database-schema-designer - Schema design and migrations
  • sqlalchemy-2-async - Async database session patterns

Capability Details

outbox-schema

Keywords: outbox table, transactional outbox, event table, schema, idempotency Solves:

  • How do I design an outbox table?
  • What indexes are needed for outbox?
  • Outbox table best practices
  • Idempotency key generation

polling-publisher

Keywords: polling, publish outbox, background worker, relay Solves:

  • How do I publish from the outbox?
  • Batch publishing patterns
  • Handling publish failures
  • FOR UPDATE SKIP LOCKED pattern

cdc-debezium

Keywords: cdc, debezium, change data capture, kafka connect Solves:

  • When to use CDC vs polling?
  • Debezium connector configuration
  • High-throughput event publishing
  • Outbox event routing

idempotent-consumer

Keywords: idempotent, deduplication, exactly-once, processed events Solves:

  • How to handle duplicate messages?
  • Idempotency key patterns
  • Consumer deduplication strategies
  • Redis + DB deduplication

atomic-operations

Keywords: atomic, transaction, dual-write, consistency Solves:

  • How to avoid dual-write problems?
  • Ensuring atomicity between DB and events
  • Exactly-once delivery patterns

适合场景

01

用户想查找某类 Agent Skill 时

02

需要根据任务场景推荐可安装能力包时

03

需要对比不同来源的安装命令和来源信息时

04

需要参考平台分布和安装热度时

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

保留来源站点、仓库和原始说明,方便继续核验

能力 4

补充不同宿主或平台的使用分布数据

能力 5

展示第三方安全扫描或审计结果

安装后应在对应宿主中按原始 README 的触发条件使用;具体调用方式请以来源页面和 README 为准。

平台分布

Claude Code

31.26%
按下载量换算31

windsurf

24.74%
按下载量换算24

trae

18.86%
按下载量换算19

OpenCode

12.82%
按下载量换算13

Codex

7.03%
按下载量换算7

Antigravity

3.42%
按下载量换算3

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

未通过

权限和风险

需要联网

该 Skill 可能需要联网访问来源站点、仓库或外部 API;具体网络访问范围需要结合源码和 README 复核。

安装前确认

本站仅展示第三方公开信息,不托管安装包,不提供自动安装或运行环境。安装前应自行审查源码、依赖和命令行为。来源安全扫描存在 warning/failed 结果,不能写成本站确认安全。

来源信息

继续浏览同类 Skills