切换主题
职责、耦合与内聚
问题:一个函数为什么越来越难改
假设报表服务的 HTTP Handler 同时读取参数、计算价格、拼 SQL、发送邮件。改邮件供应商要动它,改折扣规则也要动它,新增 CLI 更要复制它。代码不长,但多个变化来源挤在同一个地方。
架构首先解决的是变化如何传播。内聚关注一个模块里的内容是否共同完成一项责任;耦合关注模块之间知道对方多少细节。文件大小、接口数量、目录层级都不能直接衡量它们。
原理:沿着变化原因切开
把输入解析放在 HTTP 边界,把“生成报表”的步骤放在应用用例,把价格规则放在领域,把持久化放在数据库适配器。这样更换邮件服务不应改价格计算,新增 CLI 不应复制生成报表的规则。
业务边界与技术分层是两条轴。订单、计费、配送可分别有自己的应用层与适配器;把全公司所有规则塞进一个 services/ 并没有建立业务边界。反过来,每张数据库表一个模块也会拆散需要共同保护的规则。
最小例子:依赖里只放必要的信息
go
// 教学片段:省略构造和错误包装。
type ReportRequest struct {
CustomerID string
Period string
}
type ReportStore interface {
Save(context.Context, Report) error
}
func (s *ReportService) Generate(ctx context.Context, req ReportRequest) error {
report, err := BuildReport(req.CustomerID, req.Period)
if err != nil {
return err
}
return s.store.Save(ctx, report)
}ReportRequest 不带 HTTP Request;BuildReport 不接数据库连接。核心知道业务输入和协作者能力,而不是所有传输细节。用例返回错误,让 HTTP 和 CLI 各自转换成合适反馈。
交互 01
运行时调用与代码依赖
点击一层,查看它应承担的责任。两条箭头表达不同关系。
编排领域规则与事务性协作。调用核心定义的出站端口,实例由组合根注入。
- 运行时调用
- HTTP → 应用用例 → Repository 实例 → PostgreSQL
- 代码依赖
- HTTP → 核心;PostgreSQL Adapter → 核心端口
接口只是表达边界的工具。只有核心不导入基础设施实现,依赖才真正反转。这里按责任表达结构,不意味着每个小功能都需要六个文件。
方案比较
| 组织方式 | 适用情况 | 主要成本 |
|---|---|---|
| 小型垂直功能包 | 简单功能、单一业务 | 随规模增长需明确规则和外部边界 |
| 按技术层全局组织 | 业务很少且关系简单 | 同一个业务变化横跨很多目录 |
| 按业务上下文,再在内部划层 | 多组独立规则与协作关系 | 需管理跨上下文契约和共享内容 |
| 独立服务 | 部署、扩缩容或团队自治有实际需求 | 网络、数据一致性、运维成本增加 |
模块化单体可以有清晰边界。微服务也可以通过共享数据库和相互调用高度耦合。是否拆进程应由运行与组织需求决定。
真实案例:按责任组织核心
Workbench Server 把 execution、scheduling、callback 等责任分开,内部再组织 domain、application、ports。Server 和 Worker 是独立应用,通过契约通信。组合根创建数据库与消息依赖,而创建任务用例依赖 Repository、WorkflowCatalog、Clock 等能力。
已核对的实现 · 本地代码快照
创建用例与原子持久化 backend · c6dbe05c
server/internal/core/execution/application/create_task.go · 第 67–211 行
符号:CreateTask.Handle · 核对日期 2026-10-02
来源与提交版本一致
// Handle 的业务流程:校验应用、输入、回调和幂等键,解析已发布流程,
// 合并步骤调度约束并校验输入契约,最后原子写入 Task、首个 Execution、创建事件和 Outbox。
func (handler *CreateTask) Handle(ctx context.Context, command CreateTaskCommand) (*CreateTaskResult, error) {
applicationID, err := domain.ParseID(command.ApplicationID)
if err != nil {
return nil, fmt.Errorf("application ID: %w", err)
}
input, err := domain.ParseDocument(command.Input)
if err != nil {
return nil, fmt.Errorf("input: %w", err)
}
if err := validateCallbackURL(command.CallbackURL, handler.allowInsecureCallbacks); err != nil {
return nil, err
}
if strings.TrimSpace(command.PrincipalID) == "" || strings.TrimSpace(command.IdempotencyKey) == "" {
return nil, errors.New("principal ID and idempotency key are required")
}
workflow, err := handler.workflows.ResolvePublished(ctx, strings.TrimSpace(command.WorkflowName), strings.TrimSpace(command.WorkflowVersion))
if err != nil {
return nil, fmt.Errorf("resolve workflow: %w", err)
}
constraints := command.StepConstraints
if len(constraints) == 0 {
constraints = []byte(`{}`)
}
workflowSnapshot, err := mergeStepConstraints(workflow.Snapshot, constraints)
if err != nil {
return nil, fmt.Errorf("%w: %v", fault.ErrInvalid, err)
}
if handler.inputValidator != nil {
var definition struct {
InputSchema json.RawMessage `json:"input_schema"`
}
if err := json.Unmarshal(workflow.Snapshot, &definition); err != nil {
return nil, fmt.Errorf("decode workflow input contract: %w", err)
}
if err := handler.inputValidator.ValidateDocument(definition.InputSchema, command.Input, "workflow input"); err != nil {
return nil, fmt.Errorf("%w: %v", fault.ErrInvalid, err)
}
}
ids, err := generateIDs(handler.ids, 4)
if err != nil {
return nil, err
}
now := handler.clock.Now().UTC()
task, execution, err := domain.NewTask(domain.NewTaskParams{
TaskID: ids[0],
ExecutionID: ids[1],
ApplicationID: applicationID,
Name: command.Name,
ExternalRef: command.ExternalRef,
WorkflowDefinitionID: workflow.ID,
Input: input,
WorkflowSnapshot: workflowSnapshot,
StepConstraints: constraints,
CallbackURL: command.CallbackURL,
Priority: command.Priority,
CreatedAt: now,
})
if err != nil {
return nil, err
}
eventData, err := json.Marshal(map[string]any{
"task_id": string(task.ID),
"execution_id": string(execution.ID),
"external_ref": task.ExternalRef,
"status": execution.Status,
})
if err != nil {
return nil, fmt.Errorf("encode execution event: %w", err)
}
event := domain.Event{
ID: ids[2],
OutboxID: ids[3],
ApplicationID: applicationID,
ExecutionID: execution.ID,
Sequence: 1,
Type: domain.EventExecutionCreated,
Data: eventData,
OccurredAt: now,
}
requestHash := sha256.Sum256([]byte(strings.Join([]string{
command.ApplicationID, command.Name, command.ExternalRef, command.WorkflowName, command.WorkflowVersion, string(command.Input), string(constraints), command.CallbackURL, fmt.Sprint(command.Priority),
}, "\x00")))
persistedTask, persistedExecution, err := handler.repository.Create(ctx, task, execution, event, ports.Idempotency{
PrincipalID: command.PrincipalID,
Operation: "create-task",
Key: command.IdempotencyKey,
RequestHash: requestHash[:],
})
if err != nil {
return nil, fmt.Errorf("persist task: %w", err)
}
return &CreateTaskResult{
Task: persistedTask,
Execution: persistedExecution,
}, nil
}
func mergeStepConstraints(snapshot, raw []byte) ([]byte, error) {
var document map[string]any
var overrides map[string]map[string]any
if err := json.Unmarshal(snapshot, &document); err != nil {
return nil, errors.New("invalid workflow snapshot")
}
if err := json.Unmarshal(raw, &overrides); err != nil {
return nil, errors.New("invalid step constraints")
}
steps, ok := document["steps"].([]any)
if !ok {
return nil, errors.New("workflow steps are missing")
}
known := make(map[string]bool, len(steps))
for _, value := range steps {
step, ok := value.(map[string]any)
if !ok {
return nil, errors.New("invalid workflow step")
}
key, _ := step["key"].(string)
known[key] = true
override, exists := overrides[key]
if !exists {
continue
}
affinity, _ := step["affinity"].(map[string]any)
if affinity == nil {
affinity = map[string]any{}
}
for _, field := range []string{"node_id", "group"} {
if value, exists := override[field]; exists {
affinity[field] = value
}
}
for _, field := range []string{"required_labels", "preferred_labels", "minimum_resources"} {
if value, exists := override[field]; exists {
base, _ := affinity[field].(map[string]any)
if base == nil {
base = map[string]any{}
}
values, ok := value.(map[string]any)
if !ok {
return nil, fmt.Errorf("step %s %s must be an object", key, field)
}
for k, v := range values {片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
组合根及资源创建 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())
}片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
阅读时画两张图:一张是运行时调用图,一张是包导入图。一次调用跨多个模块并不自动说明耦合错误;核心必须导入具体数据库包才能工作才是需要审视的信号。
失败与边界:共享包会变成隐形大模块
把“不知道放哪”的内容都放进 shared 会让每个模块依赖共同细节。共享应有具体、稳定、狭窄的语义,例如分页或标识,而不是一套同时知道任务、用户和数据库的工具集合。
检查边界时追问:谁拥有这个规则?谁有权修改状态?别人是否绕过拥有者直接写表?跨上下文对象可以传标识和契约值,避免把整个内部实体当作公共 DTO。
迁移练习与参考答案
练习:电商下单 Handler 同时修改库存、保存订单、调用支付和发送短信。如何拆分?是否直接拆成四个微服务?
参考答案:先识别库存保留、订单状态、支付授权的业务所有者;输入和凭据校验留在边界,应用用例协调步骤,各领域保护自身规则。短信通过可靠异步通知协作。支付跨系统不属于订单数据库事务,需要幂等与对账。进程拆分仍要证明独立部署等收益,不因出现四个名词就拆服务。
判断标准:能说明每个模块因什么变化而修改,以及跨边界失败如何处理。