切换主题
六边形与依赖反转
问题:业务为什么只能通过 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 地址、认证和响应映射属于适配器;更换供应商仅改适配器与组合根。超时后是否实际发送成功不能由接口消除,需独立可靠性设计。