Skip to content

DDD 与领域建模 ​

问题:同一个“任务”为什么有不同含义 ​

在报表业务中,任务表示用户想要的报表;在执行平台中,任务表示一份可调度请求;在运行器里,任务可能是一个进程。把这三种模型直接共享,会让业务字段、调度细节和进程状态互相污染。

DDD 首先要求在具体上下文中使用一致语言。限界上下文规定一个模型与语言的适用范围,并明确其他上下文如何与它协作。它不是每个包、每张表或每个微服务的别名。

原理:从身份与不变量建模 ​

实体凭身份延续,即使名称改变还是同一对象。值对象凭内容表达含义,例如金额由数值与币种组成,不允许负数时就在构造入口保护规则。领域模型不是“数据库结构体加 getter”。

聚合是一组需要共同维护一致性的对象及其入口。聚合根保护不变量,事务边界通常围绕需要原子成立的规则设计。对象引用很多并不意味着全部必须归进一个大聚合。

应用服务编排用例:读取、授权、调用领域规则、保存。领域服务承载不自然属于单个实体的业务规则。发送 HTTP、写 SQL 等技术协作属于适配器。

最小例子:构造有效值,保护行为 ​

go
// 教学片段:Amount 使用分,不使用浮点数表达金额。
type Money struct {
    cents    int64
    currency string
}

func NewMoney(cents int64, currency string) (Money, error) {
    if cents < 0 || currency == "" {
        return Money{}, errors.New("invalid money")
    }
    return Money{cents: cents, currency: currency}, nil
}

func (order *Order) Cancel() error {
    if order.status == "SHIPPED" { return errors.New("already shipped") }
    order.status = "CANCELED"
    return nil
}

此例只展示局部规则,保存仍要处理并发。两次请求都读到未发货并不意味着取消一定有效;事务中应基于版本或当前状态条件确认修改。

可下载这一最小示例的 Go 源码、测试与 go.mod,放在同一目录运行 go test ./...。测试覆盖无效金额、零金额、已发货拒绝取消及重复取消;示例没有数据库并发保证。

方案比较:模型按复杂度生长 ​

方式优势应注意
事务脚本CRUD 简单直接重复规则容易散落入口
领域模型集中状态规则与不变量不必给每个字段发明值对象
富跨表聚合强事务一致性大范围锁与保存成本可能过高
小聚合与事件协作独立维护与扩展需要明确中间状态和补偿

DDD 与六边形处理不同问题:前者帮助建立业务模型,后者保护模型与技术的边界。采用目录格式不等于完成了领域设计。

真实案例:长期意图与一次运行 ​

TDP 用 Task 表达长期请求,用 Execution 表达一次运行。Execution 冻结流程快照,重跑产生新运行。领域的 ID 是强类型字符串,Document 在解析时校验并复制 JSON 字节,避免调用者复用缓冲区后改变内部值。

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

实体、值对象与创建不变量 backend · c6dbe05c

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

// ID 是 Execution 上下文使用的强类型资源标识,避免把普通字符串误传为领域 ID。
type ID string

// ParseID 去除首尾空白并拒绝空标识。
func ParseID(value string) (ID, error) {
	value = strings.TrimSpace(value)
	if value == "" {
		return "", errors.New("ID is required")
	}
	return ID(value), nil
}

// Document 保存已经校验过的 JSON 文档,并在进入领域层时复制以避免外部修改。
type Document []byte

// ParseDocument 校验并复制 JSON 字节,调用方可安全复用原缓冲区。
func ParseDocument(value []byte) (Document, error) {
	if len(value) == 0 || !json.Valid(value) {
		return nil, errors.New("document must contain valid JSON")
	}
	return append(Document(nil), value...), nil
}

// Task 是业务请求的长期聚合根;一次 Task 可以通过重跑产生多个 Execution。
type Task struct {
	ID                   ID              // 平台生成的任务标识。
	ApplicationID        ID              // 任务的归属和访问范围。
	Name                 string          // 供业务人员识别的任务名称。
	ExternalRef          string          // 业务系统用于关联和幂等追踪的外部引用。
	WorkflowDefinitionID ID              // 创建任务时解析到的已发布流程。
	Input                Document        // 业务系统提交给流程的原始输入。
	StepConstraints      Document        // 按步骤合并后的调度约束。
	EffectiveStatus      ExecutionStatus // 当前执行状态,便于直接查询任务。
	CurrentExecutionID   ID              // 最近一次创建或重跑的执行。
	Priority             int16           // 待调度任务的业务优先级,范围为 -100 到 100。
	Version              int64           // 用于任务的乐观并发控制。
	CreatedAt            time.Time       // 任务首次创建时间。
	UpdatedAt            time.Time       // 聚合最后一次状态变更时间。
}

// Execution 是 Task 的一次具体运行,保存创建时冻结的流程快照。
type Execution struct {
	ReconciliationError string          // 终态对账失败或待人工检查的原因。
	SchedulingReasons   []string        // 当前排队步骤的等待原因。
	ID                  ID              // 本次执行的稳定标识。
	TaskID              ID              // 所属任务。
	Number              int32           // 同一任务下从 1 开始递增的执行序号。
	Status              ExecutionStatus // 执行生命周期状态。
	Progress            int16           // 0 到 100 的单调进度值。
	Result              Document        // 执行完成后的 JSON 结果。
	WorkflowSnapshot    Document        // 冻结本次运行使用的已编译流程。
	CallbackURL         string          // 完成事件的业务回调地址。
	Reason              string          // 记录重跑、取消或失败等人工可读原因。
	Version             int64           // 用于执行的乐观并发控制。
	CreatedAt           time.Time       // 执行创建时间。
	StartedAt           *time.Time      // 首次进入运行态的时间,未开始时为空。
	CompletedAt         *time.Time      // 进入终态的时间,未结束时为空。
}

// NewTaskParams 汇集原子创建 Task 及首个 Execution 所需的已解析参数。
type NewTaskParams struct {
	TaskID               ID        // 预先生成的任务标识。
	ExecutionID          ID        // 首个执行标识。
	ApplicationID        ID        // 任务所属应用。
	Name                 string    // 任务展示名称。
	ExternalRef          string    // 业务系统外部引用。
	WorkflowDefinitionID ID        // 已发布流程标识。
	Input                Document  // 已经校验的任务输入。
	WorkflowSnapshot     Document  // 创建时冻结的流程定义。
	StepConstraints      Document  // 已经合并的步骤调度约束。
	CallbackURL          string    // 可选的执行事件回调地址。
	Priority             int16     // 任务优先级。
	CreatedAt            time.Time // 统一注入的创建时间。
}

// NewTask 校验创建不变量,并同时构造 Task 与编号为 1 的 Execution。
func NewTask(params NewTaskParams) (*Task, *Execution, error) {
	if err := validateIDs(params.TaskID, params.ExecutionID, params.ApplicationID, params.WorkflowDefinitionID); err != nil {
		return nil, nil, err
	}
	if strings.TrimSpace(params.Name) == "" {
		return nil, nil, errors.New("task name is required")
	}
	if len([]rune(strings.TrimSpace(params.Name))) > 200 {
		return nil, nil, errors.New("task name must not exceed 200 characters")
	}
	if strings.TrimSpace(params.ExternalRef) == "" {
		return nil, nil, errors.New("external reference is required")
	}
	if params.Priority < -100 || params.Priority > 100 {
		return nil, nil, errors.New("priority must be between -100 and 100")
	}
	if len(params.StepConstraints) == 0 {
		params.StepConstraints = Document(`{}`)
	}
	if !json.Valid(params.Input) || !json.Valid(params.WorkflowSnapshot) || !json.Valid(params.StepConstraints) {
		return nil, nil, errors.New("input and workflow snapshot must contain valid JSON")
	}
	if params.CreatedAt.IsZero() {
		return nil, nil, errors.New("creation time is required")
	}

	execution := &Execution{
		ID:               params.ExecutionID,
		TaskID:           params.TaskID,
		Number:           1,
		Status:           ExecutionCreated,
		WorkflowSnapshot: append(Document(nil), params.WorkflowSnapshot...),
		CallbackURL:      params.CallbackURL,
		Version:          1,
		CreatedAt:        params.CreatedAt,
	}
	task := &Task{
		ID:                   params.TaskID,
		ApplicationID:        params.ApplicationID,
		Name:                 strings.TrimSpace(params.Name),
		ExternalRef:          strings.TrimSpace(params.ExternalRef),
		WorkflowDefinitionID: params.WorkflowDefinitionID,
		Input:                append(Document(nil), params.Input...),
		StepConstraints:      append(Document(nil), params.StepConstraints...),
		EffectiveStatus:      execution.Status,
		CurrentExecutionID:   execution.ID,
		Priority:             params.Priority,
		Version:              1,
		CreatedAt:            params.CreatedAt,
		UpdatedAt:            params.CreatedAt,
	}
	return task, execution, nil
}

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

业务创建与命令提交边界 app-demo · c383860e

internal/core/task/application/service.go · 第 50–111 行
符号:Service.Create · 核对日期 2026-10-02
来源与提交版本一致

func (s *Service) Create(ctx context.Context, name, sourceAssetID string, processingMode taskdomain.ProcessingMode, watermark, idempotencyKey string) (*taskdomain.Order, error) {
	name, watermark = strings.TrimSpace(name), strings.TrimSpace(watermark)
	if name == "" || len(name) > 200 {
		return nil, errors.New("name must contain 1 to 200 characters")
	}
	if !processingMode.Valid() {
		return nil, errors.New("processing_mode must be ARCHIVE_CLEAN_ONLY or ARCHIVE_CLEAN_THEN_WATERMARK")
	}
	if processingMode == taskdomain.ProcessingModeArchiveCleanThenWatermark && (watermark == "" || len([]rune(watermark)) > 100) {
		return nil, errors.New("watermark_text must contain 1 to 100 characters")
	}
	if processingMode == taskdomain.ProcessingModeArchiveCleanOnly && watermark != "" {
		return nil, errors.New("watermark_text must be empty when processing_mode is ARCHIVE_CLEAN_ONLY")
	}
	asset, err := s.repository.GetAsset(ctx, sourceAssetID)
	if err != nil {
		return nil, err
	}
	if asset.Kind != assetdomain.KindSourceZIP || asset.State != "READY" {
		return nil, errors.New("source asset is not a ready ZIP")
	}
	requestHash := hashRequest(map[string]string{"name": name, "processing_mode": string(processingMode), "source_asset_id": sourceAssetID, "watermark_text": watermark})
	if replay, found, replayErr := s.replay(ctx, "CREATE", idempotencyKey, requestHash); found || replayErr != nil {
		return replay, replayErr
	}
	id, now := uuid.NewString(), time.Now().UTC()
	order, err := s.repository.CreateOrder(ctx, taskdomain.Order{
		ID:                  id,
		Name:                name,
		SourceAssetID:       sourceAssetID,
		OutputCollectionRef: assetdomain.CollectionReference(id),
		ProcessingMode:      processingMode,
		WatermarkText:       watermark,
		BusinessStatus:      taskdomain.StatusSubmitting,
		Result:              []byte(`{}`),
		CreatedAt:           now,
		UpdatedAt:           now,
	}, idempotencyKey, requestHash)
	if err == nil {
		return order, nil
	}
	if replay, found, replayErr := s.replay(ctx, "CREATE", idempotencyKey, requestHash); found || replayErr != nil {
		return replay, replayErr
	}
	return nil, err
}

func (s *Service) Get(ctx context.Context, id string) (*taskdomain.Order, error) {
	return s.repository.GetOrder(ctx, id)
}
func (s *Service) List(ctx context.Context) ([]taskdomain.Order, int64, error) {
	return s.repository.ListOrders(ctx, 100, 0)
}
func (s *Service) Cancel(ctx context.Context, id, key string) (*taskdomain.Order, error) {
	requestHash := hashRequest(map[string]string{"operation": "cancel", "order_id": id})
	if replay, found, err := s.replay(ctx, "CANCEL", key, requestHash); found || err != nil {
		return replay, err
	}
	current, err := s.repository.GetOrder(ctx, id)
	if err != nil {
		return nil, err
	}

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

App Demo 的处理单还拥有资源、文字水印和业务状态,这是业务上下文的模型。TDP 只接收资源引用与流程输入,不需要知道某个业务为什么要加水印。

注意:不同领域类型处于同一上下文,并不自动意味着它们都在同一个聚合事务里。确定边界要继续读写入用例与数据库事务。

失败与边界:避免强行套术语 ​

公开字段的 Go 实体并不能完全阻止所有调用者直接赋值;构造与行为方法、应用流程和持久化条件需要共同保护规则。共享 DTO 更不能直接作为各上下文共同的领域模型。

领域事件描述已经发生的事实,集成消息是向外部传递事实的契约,两者可能不同。将领域对象直接 JSON 序列化对外会把内部变化变成接口破坏。

迁移练习与参考答案 ​

练习:视频平台的“转码任务”和算力平台的“执行任务”是否共用一套结构体?用户重试后应覆盖原来运行吗?

参考答案:分别拥有业务模型与执行模型,以稳定业务引用和执行标识关联。视频平台维护用户、源视频和发布规则;算力平台维护步骤、尝试与运行状态。重试应保留上次运行及失败证据,创建新执行或新尝试,区别业务重跑与内部重试。

继续阅读:任务模型与状态机、业务接入与投影。

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