Skip to content

事务与一致性边界 ​

问题:成功响应究竟承诺了什么 ​

用户创建报表任务后立即关闭页面。接口已经返回成功,但进程还没有把任务放进队列。如果队列动作只保存在内存,成功响应承诺的工作就可能消失。

设计事务时先写一句承诺:“创建成功意味着请求、首个执行和推进所需的待发事实均已持久化。”然后推导必须同事务写入的记录,而不是先讨论某个 ORM 的 Begin 方法。

原理:原子性与并发正确性是两件事 ​

事务让一组数据库修改一起提交或回滚。隔离级别决定并发事务能看到什么;一个事务不自动避免丢失更新、超卖或不正确的状态判断。

例如先读库存为 1,再无条件写成 0,两个事务可能都认为自己买到最后一件。可以条件更新 stock >= quantity 并检查受影响行数,或在事务里锁住行。唯一约束则保护某类事实至多一条。

最小例子:条件更新保护不变量 ​

sql
-- 教学 SQL:失败时应用必须回滚整个事务。
BEGIN;
UPDATE products
SET stock = stock - 1
WHERE id = $1 AND stock >= 1;
-- 必须检查受影响行数为 1,否则不能创建订单。
INSERT INTO orders (id, product_id, status)
VALUES ($2, $1, 'PLACED');
INSERT INTO outbox (id, aggregate_id, type)
VALUES ($3, $2, 'order.placed');
COMMIT;

订单与 Outbox 在这里属于同一数据库事务;消息发送在提交后完成。消息系统不参加此数据库事务,因此事务成功不能保证此刻接收端已经处理消息。

方案比较 ​

手段保护什么局限
行锁事务内读取和修改同一行竞争、等待、死锁需要管理
乐观版本拒绝基于旧快照的更新冲突后需重新决策,不应盲目覆盖
唯一约束同一业务键只能有一条记录不能替代复杂状态规则
Serializable更强的并发事务隔离可能中止事务,需要安全重试
Outbox 与补偿跨系统逐步推进存在中间状态,需幂等与观测

长时间远端调用不宜放在持锁事务里。网络延迟扩大锁时间,而且远端成功后数据库回滚仍不能撤回远端效果。把跨系统业务建模为可恢复步骤,单独记录中间状态。

真实案例:端口承诺,适配器兑现 ​

创建 Task 用例构造 Task、Execution、事件和幂等参数,Repository.Create 的契约承诺原子持久化。Attempt 调度事务同时检查执行状态、调度认领和重复 Attempt,并继续预留 Worker、写命令。

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

核心定义协作契约 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/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 {

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

调度事务与幂等 Attempt backend · c6dbe05c

server/internal/adapters/outbound/postgres/attempt_dispatch.go · 第 45–133 行
符号:dispatchAttempt · 核对日期 2026-10-02
来源与提交版本一致

func (dispatcher *AttemptDispatcher) dispatchAttempt(ctx context.Context, request executionports.CreateAttemptRequest, claim *pgqueries.SchedulingRequest) (*executionports.AttemptRef, error) {
	stepID, err := uuid.Parse(request.StepID)
	if err != nil {
		return nil, fmt.Errorf("parse step ID: %w", err)
	}
	tx, err := dispatcher.database.BeginTx(ctx, pgx.TxOptions{
		IsoLevel: pgx.ReadCommitted,
	})
	if err != nil {
		return nil, err
	}
	defer func() { _ = tx.Rollback(ctx) }()
	txQueries := dispatcher.queries.WithTx(tx)
	lockKey := request.StepID + ":" + fmt.Sprint(request.AttemptNumber)
	if err := txQueries.AcquireTransactionAdvisoryLock(ctx, lockKey); err != nil {
		return nil, err
	}
	if _, err := txQueries.LockSchedulingRequest(ctx, pgqueries.LockSchedulingRequestParams{
		ID:      claim.ID,
		ClaimID: claim.ClaimID,
	}); err != nil {
		return nil, err
	}
	executionIDValue, err := uuid.Parse(request.ExecutionID)
	if err != nil {
		return nil, fmt.Errorf("%w: execution ID", executiondomain.ErrDispatchInvalid)
	}
	execution, err := txQueries.LockExecutionByID(ctx, executionIDValue)
	if err != nil {
		return nil, err
	}
	if execution.Status != "CREATED" && execution.Status != "QUEUED" && execution.Status != "RUNNING" {
		return nil, fmt.Errorf("%w: execution closed", executiondomain.ErrDispatchInvalid)
	}
	existing, err := txQueries.FindAttemptByStepAndNumber(ctx, pgqueries.FindAttemptByStepAndNumberParams{
		StepID:        stepID,
		AttemptNumber: int32(request.AttemptNumber),
	})
	if err == nil {
		if !existing.WorkerID.Valid {
			return nil, errors.New("existing attempt has no worker")
		}
		if err := txQueries.CompleteSchedulingRequest(ctx, pgqueries.CompleteSchedulingRequestParams{
			ID:          claim.ID,
			ClaimID:     claim.ClaimID,
			State:       "DISPATCHED",
			AvailableAt: timestamp(time.Now().UTC()),
		}); err != nil {
			return nil, err
		}
		if err := tx.Commit(ctx); err != nil {
			return nil, err
		}
		existingWorker := uuid.UUID(existing.WorkerID.Bytes)
		_ = dispatcher.waker.WakeWorker(existingWorker.String())
		return &executionports.AttemptRef{
			AttemptID:    existing.ID.String(),
			LeaseVersion: existing.LeaseVersion,
		}, nil
	}
	if !errors.Is(err, pgx.ErrNoRows) {
		return nil, err
	}
	now := time.Now().UTC()
	workerID, generation, err := dispatcher.reserveWorker(ctx, txQueries, stepID, request, now)
	if err != nil {
		if errors.Is(err, errNoEligibleWorker) {
			if commitErr := tx.Commit(ctx); commitErr != nil {
				return nil, commitErr
			}
		}
		return nil, err
	}
	contextRow, err := txQueries.FindStepExecutionAndWorkerNode(ctx, pgqueries.FindStepExecutionAndWorkerNodeParams{
		StepID:   stepID,
		WorkerID: workerID,
	})
	if err != nil {
		return nil, err
	}
	executionID, nodeID := contextRow.ExecutionID, contextRow.NodeID
	if err := appendTimelineTx(ctx, txQueries, timelineRecord{
		ExecutionID:   executionID,
		MilestoneKey:  "step.worker.selected:" + stepID.String() + ":" + fmt.Sprint(request.AttemptNumber),
		Stage:         "SCHEDULING",
		AttemptNumber: request.AttemptNumber,
		State:         "SUCCEEDED",
		Component:     "Scheduler",
		Title:         "已选定 Node 与 Capability 唯一 Worker",

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

领域对象负责表达规则,持久化适配器负责并发和原子写入。不要从 NewTask 同时构造两个对象就推出数据库事务已经正确;需要追到 Repository 实现。

失败与边界:提交结果未知 ​

客户端超时时,数据库可能已经提交,只是响应丢失。不能把“没收到成功”解释为“肯定失败”。应使用稳定业务标识与幂等请求键再次查询或重试。

事务提交前不发送成功响应。辅助日志与审计的持久化要求也要明确:哪些失败必须阻止业务提交,哪些只记录告警。所有写入无限塞进一个大事务会掩盖这种责任区别。

迁移练习与参考答案 ​

练习:创建订单时要扣本地库存、保存订单、向外部支付服务预授权。怎样安排事务?

参考答案:本地保留库存与订单待支付状态、支付请求意图在同一事务中提交;持久化步骤执行外部预授权,用稳定支付请求 ID 防重复。将支付结果在新事务中推进订单;失败时释放库存,未知时对账。不能持锁调用支付并假定回滚能够撤销授权。

继续阅读:Outbox、锁与租约。

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