切换主题
事务与一致性边界
问题:成功响应究竟承诺了什么
用户创建报表任务后立即关闭页面。接口已经返回成功,但进程还没有把任务放进队列。如果队列动作只保存在内存,成功响应承诺的工作就可能消失。
设计事务时先写一句承诺:“创建成功意味着请求、首个执行和推进所需的待发事实均已持久化。”然后推导必须同事务写入的记录,而不是先讨论某个 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 防重复。将支付结果在新事务中推进订单;失败时释放库存,未知时对账。不能持锁调用支付并假定回滚能够撤销授权。