Skip to content

六边形与依赖反转 ​

问题:业务为什么只能通过 HTTP 运行 ​

一个“审批申请”规则如果必须拿到 Web 框架 Context 和 ORM Model 才能调用,那么测试它需要搭建整个系统,增加消息入口也要模仿 HTTP。这说明业务与交付技术黏在了一起。

六边形架构把应用视为核心,通过端口与外部交互。六条边没有数量约束;重要的是核心逻辑可以由不同入口驱动,也可以通过不同技术完成外部协作。

原理:代码依赖向核心 ​

入站端口表达应用可以做什么,例如 ApproveApplication;入站适配器把 HTTP、CLI 或消息转成这个命令。出站端口表达核心需要什么,例如查申请、保存审批结果、获取时间。

运行时用例会调用数据库适配器;代码上适配器实现核心定义的接口。调用方向与依赖方向不同,这是依赖反转的要点。端口方向由谁发起协作决定,不按网络包是收还是发来命名。

最小例子:业务定义接口,外部实现 ​

go
// 核心契约:教学片段。
type ApplicationRepository interface {
    Load(context.Context, string) (Application, error)
    SaveApproved(context.Context, Application) error
}

type Approver struct { repository ApplicationRepository }

func (a *Approver) Approve(ctx context.Context, id string) error {
    application, err := a.repository.Load(ctx, id)
    if err != nil { return err }
    if err := application.Approve(); err != nil { return err }
    return a.repository.SaveApproved(ctx, application)
}

PostgreSQL Adapter 实现两个方法,组合根将实例传给 Approver。真实并发审批还需版本校验或锁,接口本身不提供事务保证。端口应描述这一业务所需的原子保存语义。

交互 01

运行时调用与代码依赖

点击一层,查看它应承担的责任。两条箭头表达不同关系。

编排领域规则与事务性协作。调用核心定义的出站端口,实例由组合根注入。
运行时调用
HTTP → 应用用例 → Repository 实例 → PostgreSQL
代码依赖
HTTP → 核心;PostgreSQL Adapter → 核心端口

接口只是表达边界的工具。只有核心不导入基础设施实现,依赖才真正反转。这里按责任表达结构,不意味着每个小功能都需要六个文件。

方案比较:抽象必须有意义 ​

直接使用具体数据库类在简单脚本里成本很低;长期业务规则需要独立测试和替换基础设施时,端口能保护边界。把 Query(sql string, args ...any) 包装成接口并没有隔离数据库知识,核心仍在拼 SQL。

接口过大也会带来耦合:只需要读订单的调用者不该同时依赖管理员删除、导出、迁移方法。以协作责任拆端口,而不是每个函数都制造一个接口。

Go 接口可由消费者定义;在本案例中集中于 core/ports,便于表达边界和架构检查。这是工程组织选择,不是所有 Go 项目的唯一目录格式。

真实案例:语义化持久化与组合根 ​

Server 的 Repository.Create 明确承担 Task、Execution、事件和幂等记录的原子持久化。核心传入领域对象与请求指纹,不传 pgx 事务。OutboxDispatcher 只依赖 OutboxStore、EventPublisher 和 Clock。

已核对的实现 · 本地代码快照

核心定义协作契约 backend · c6dbe05c

server/internal/core/execution/ports/outbound/repository.go · 第 9–24 行
符号:Repository · 核对日期 2026-10-02
来源与提交版本一致

// Idempotency 描述一次写操作的调用方、键和请求指纹。
type Idempotency struct {
	PrincipalID string // 隔离不同调用方的同名幂等键。
	Operation   string // 隔离同一调用方的不同业务动作。
	Key         string // 调用方提供的幂等键。
	RequestHash []byte // 用于拒绝复用同一键但内容不同的请求。
}

// Repository 原子持久化 Task、Execution、首个事件和幂等记录。
type Repository interface {
	// Create 创建聚合;重复且请求指纹一致时返回原结果。
	Create(context.Context, *domain.Task, *domain.Execution, domain.Event, Idempotency) (*domain.Task, *domain.Execution, error)
	// Get 在应用范围内读取任务。
	Get(context.Context, Scope, domain.ID) (*domain.Task, error)
}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

发布、认领与退避 backend · c6dbe05c

server/internal/core/execution/application/outbox_dispatcher.go · 第 12–81 行
符号:DispatchOutboxBatch · 核对日期 2026-10-02
来源与提交版本一致

const (
	outboxClaimTTL    = time.Minute
	outboxMaxAttempts = 20
)

// OutboxDispatcher 将领域事务中的 Outbox 事件可靠发布到消息系统。
type OutboxDispatcher struct {
	store     executionports.OutboxStore
	publisher executionports.EventPublisher
	clock     executionports.Clock
}

func NewOutboxDispatcher(store executionports.OutboxStore, publisher executionports.EventPublisher, clock executionports.Clock) (*OutboxDispatcher, error) {
	if store == nil || publisher == nil || clock == nil {
		return nil, errors.New("outbox store, event publisher, and clock are required")
	}
	return &OutboxDispatcher{
		store:     store,
		publisher: publisher,
		clock:     clock,
	}, nil
}

func (dispatcher *OutboxDispatcher) DispatchOutboxBatch(ctx context.Context, limit int) (int, []error) {
	if limit < 1 {
		return 0, []error{errors.New("outbox batch limit must be positive")}
	}
	now := dispatcher.clock.Now().UTC()
	events, err := dispatcher.store.ClaimOutboxEvents(ctx, now, outboxClaimTTL, limit)
	if err != nil {
		return 0, []error{fmt.Errorf("claim outbox events: %w", err)}
	}
	failures := make([]error, 0)
	for _, event := range events {
		err := dispatcher.publisher.PublishExecutionEvent(ctx, executionports.PublishedEvent{
			ID:            event.ID,
			Type:          event.EventType,
			CorrelationID: event.AggregateID,
			Sequence:      event.Sequence,
			OccurredAt:    event.CreatedAt,
			Payload:       event.Payload,
		})
		if err == nil {
			err = dispatcher.store.MarkOutboxEventPublished(ctx, event.ID, event.ClaimID, dispatcher.clock.Now().UTC())
			if err != nil {
				failures = append(failures, fmt.Errorf("mark outbox event %s published: %w", event.ID, err))
			}
			continue
		}
		status := "FAILED"
		if event.Attempts >= outboxMaxAttempts {
			status = "DEAD"
		}
		availableAt := dispatcher.clock.Now().UTC().Add(time.Second * time.Duration(1<<min(event.Attempts, 10)))
		if markErr := dispatcher.store.MarkOutboxEventFailed(ctx, event.ID, event.ClaimID, status, availableAt, err.Error()); markErr != nil {
			err = errors.Join(err, markErr)
		}
		failures = append(failures, fmt.Errorf("publish outbox event %s: %w", event.ID, err))
	}
	return len(events), failures
}

func (dispatcher *OutboxDispatcher) NextOutboxAttemptAt(ctx context.Context) (time.Time, bool, error) {
	next, ok, err := dispatcher.store.NextOutboxAttemptAt(ctx)
	if err != nil {
		return time.Time{}, false, fmt.Errorf("read next outbox attempt: %w", err)
	}
	return next, ok, nil
}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

组合根及资源创建 backend · c6dbe05c

server/internal/bootstrap/runtime.go · 第 68–167 行
符号:run · 核对日期 2026-10-02
来源与提交版本一致

func run(ctx context.Context, serverConfig config.Server, logger *slog.Logger) (returnedErr error) {
	if logger == nil {
		return errors.New("server logger is required")
	}
	ctx, cancel := context.WithCancel(ctx)
	defer cancel()
	var startupCtx context.Context
	var startupCancel context.CancelFunc
	ctx = operation.Start(ctx, operation.Operation{
		Name: "server.run",
	})
	defer func() { returnedErr = operation.WithFailure(ctx, returnedErr) }()
	startupCtx, startupCancel = context.WithTimeout(ctx, 30*time.Second)
	defer startupCancel()
	poolConfig, err := pgxpool.ParseConfig(serverConfig.DatabaseURL)
	if err != nil {
		return fault.Wrap(fault.Invalid, "postgresql.configure", err)
	}
	poolConfig.ConnConfig.Tracer = observability.NewPostgreSQLQueryLogger(logger.With("component", "postgresql"), serverConfig.SlowQueryThreshold)
	pool, err := pgxpool.NewWithConfig(startupCtx, poolConfig)
	if err != nil {
		return fault.Wrap(fault.Unavailable, "postgresql.connect", err)
	}
	defer pool.Close()
	if err := pool.Ping(startupCtx); err != nil {
		return fault.Wrap(fault.Unavailable, "postgresql.ping", err)
	}
	var baseline string
	if err := pool.QueryRow(startupCtx, "SELECT version FROM tdp_schema_baseline").Scan(&baseline); err != nil || baseline != "application-v1" {
		return errors.New("TDP application-v1 database required; initialize a fresh project database")
	}
	bus, err := natsmessaging.Connect(serverConfig.NATSURL, "tdp-server", 10*time.Second, logger.With("component", "nats"))
	if err != nil {
		return fault.Wrap(fault.Unavailable, "nats.connect", err)
	}
	defer func() {
		if err := bus.Close(); err != nil {
			logger.ErrorContext(ctx, "close NATS", "error", err)
		}
	}()
	if err := bus.EnsureStreams(startupCtx); err != nil {
		return fault.Wrap(fault.Unavailable, "nats.ensure_streams", err)
	}
	cryptoBox, err := cryptobox.NewBase64(serverConfig.DataEncryptionKey)
	if err != nil {
		return err
	}
	definitionValidator, err := definitionvalidation.New()
	if err != nil {
		return err
	}
	grantOptions := []accessgrantapp.Option{}
	if serverConfig.AllowInsecureResourceURLs {
		logger.WarnContext(ctx, "insecure HTTP resource URLs are enabled; use only in local E2E")
		grantOptions = append(grantOptions, accessgrantapp.WithInsecureHTTPForLocalDevelopment())
	}
	grantRepository, err := accessgrantpostgres.NewRepository(pool)
	if err != nil {
		return err
	}
	resourceProvider, err := resourceprovider.New(15 * time.Second)
	if err != nil {
		return err
	}
	grantModule, err := accessgrantapp.NewModule(accessgrantapp.ModuleDependencies{
		Repository: grantRepository,
		Provider:   resourceProvider,
		Cipher:     cryptoBox,
		Clock:      clock.System{},
		IDs:        idgen.StringUUIDv7{},
		Options:    grantOptions,
	})
	if err != nil {
		return err
	}
	grantIssuer, err := workercontroladapter.NewAccessGrantIssuer(grantModule.Service)
	if err != nil {
		return err
	}
	workerStore, err := postgresadapter.NewWorkerControlStore(pool)
	if err != nil {
		return err
	}
	workerControlModule, err := workercontrol.NewModule(workercontrol.ModuleDependencies{
		Store:               workerStore,
		IDs:                 idgen.StringUUIDv7{},
		HeartbeatInterval:   serverConfig.WorkerControl.HeartbeatInterval,
		LeaseDuration:       serverConfig.WorkerControl.LeaseDuration,
		CommandPollInterval: serverConfig.WorkerControl.CommandPollInterval,
		WakeSubscriber:      bus,
		GrantIssuer:         grantIssuer,
		CapabilityValidator: definitionValidator,
	})
	if err != nil {
		return err
	}
	taskRepositoryOptions := []postgresadapter.TaskRepositoryOption{}
	if serverConfig.AllowInsecureCallbackURLs {
		taskRepositoryOptions = append(taskRepositoryOptions, postgresadapter.WithInsecureCallbacksForLocalDevelopment())
	}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

这说明端口可以规定原子效果,具体 SQL 与事务由适配器实现;但评审仍必须读适配器和测试,确认承诺被兑现。构造函数校验依赖,使配置遗漏在启动时显现。

失败与边界:绕过组合根的隐式依赖 ​

全局数据库变量、运行时自行寻找配置、服务定位器让依赖藏起来。单元测试可能意外连接真实服务。将资源创建集中在组合根,核心只使用构造时传入的协作者,可以让启动和关闭顺序可见。

领域不应把外部超时当作无效业务输入;适配器要保留可分类的错误,应用决定重试与补偿,入口决定响应格式。

迁移练习与参考答案 ​

练习:通知系统从 SMTP 切换到 HTTP 邮件服务。核心端口设计成 SendEmail(recipient, subject, body) 还是 POST(url, headers, json)?

参考答案:核心需要发送通知的语义,应使用邮件或通知契约,传稳定通知 ID 以支持幂等。如果不同渠道的业务需求不同,不把它们硬塞成万能发送器。HTTP 地址、认证和响应映射属于适配器;更换供应商仅改适配器与组合根。超时后是否实际发送成功不能由接口消除,需独立可靠性设计。

继续阅读:DDD 与领域建模、设计原则与模式。

理解原理 · 分析取舍 · 用真实代码检验