Skip to content

职责、耦合与内聚 ​

问题:一个函数为什么越来越难改 ​

假设报表服务的 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 同时修改库存、保存订单、调用支付和发送短信。如何拆分?是否直接拆成四个微服务?

参考答案:先识别库存保留、订单状态、支付授权的业务所有者;输入和凭据校验留在边界,应用用例协调步骤,各领域保护自身规则。短信通过可靠异步通知协作。支付跨系统不属于订单数据库事务,需要幂等与对账。进程拆分仍要证明独立部署等收益,不因出现四个名词就拆服务。

判断标准:能说明每个模块因什么变化而修改,以及跨边界失败如何处理。

继续阅读:六边形与依赖反转、事务与一致性边界。

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