event

package
v1.1.64 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 22, 2026 License: Apache-2.0 Imports: 4 Imported by: 0

README

Domain Event Package

概述

domain/event 包提供了企业级事件驱动架构的核心领域事件结构和工具函数。这是 jxt-core 的 DDD 共享内核(Shared Kernel)的一部分,为所有微服务提供统一的事件定义和处理标准。

核心组件

1. BaseDomainEvent(基础领域事件)

包含所有事件驱动系统都需要的核心字段,适用于所有使用事件驱动架构的系统。

核心字段:

  • EventID: 事件唯一标识(UUIDv7,保证时序性)
  • EventType: 事件类型
  • OccurredAt: 事件发生时间
  • Version: 事件版本
  • AggregateID: 聚合根ID
  • AggregateType: 聚合根类型
  • Payload: 事件载荷

使用示例:

import (
    jxtevent "github.com/ChenBigdata421/jxt-core/sdk/pkg/domain/event"
)

// 创建基础事件
event := jxtevent.NewBaseDomainEvent(
    "Archive.Created",
    "archive-123",
    "Archive",
    map[string]interface{}{
        "title": "Test Archive",
    },
)
2. EnterpriseDomainEvent(企业级领域事件)

在 BaseDomainEvent 基础上增加企业级通用字段,适用于多租户 SaaS 系统和企业级应用。

额外字段:

  • TenantId: 租户ID(类型:int,默认为 0 表示全局/无租户)
  • CorrelationId: 业务关联ID(用于业务流程追踪)
  • CausationId: 因果事件ID(用于事件因果链分析)
  • TraceId: 分布式追踪ID(集成分布式追踪系统)

TenantId 类型说明:

  • 0: 表示系统级/无租户事件
  • 1, 2, 3, ...: 表示具体租户ID
  • 与租户中间件类型一致,避免类型转换

使用示例:

// 创建企业级事件
event := jxtevent.NewEnterpriseDomainEvent(
    "Archive.Created",
    "archive-123",
    "Archive",
    payload,
)

// 设置租户ID
event.SetTenantId(1)

// 设置可观测性字段
event.SetCorrelationId("workflow-123")
event.SetCausationId("trigger-event-456")
event.SetTraceId("trace-789")
3. 序列化/反序列化助手
3.1 DomainEvent 序列化/反序列化

用于序列化和反序列化完整的 DomainEvent 对象。

使用场景:

  • Command Side: 保存到 Outbox 表
  • Query Side: 从消息队列接收事件

方法列表:

  • MarshalDomainEvent(event BaseEvent) ([]byte, error) - 序列化完整的 DomainEvent
  • UnmarshalDomainEvent[T BaseEvent](data []byte) (T, error) - 反序列化完整的 DomainEvent
  • MarshalDomainEventToString(event BaseEvent) (string, error) - 序列化为 JSON 字符串
  • UnmarshalDomainEventFromString[T BaseEvent](jsonString string) (T, error) - 从 JSON 字符串反序列化

使用示例:

// Command Side: 序列化事件保存到 Outbox
event := jxtevent.NewEnterpriseDomainEvent("Archive.Created", "archive-123", "Archive", payload)
eventBytes, err := jxtevent.MarshalDomainEvent(event)
if err != nil {
    return fmt.Errorf("failed to marshal event: %w", err)
}
// 保存 eventBytes 到 Outbox 表

// Query Side: 从消息队列接收并反序列化
receivedEvent, err := jxtevent.UnmarshalDomainEvent[*jxtevent.EnterpriseDomainEvent](msg.Payload)
if err != nil {
    return fmt.Errorf("failed to unmarshal event: %w", err)
}
3.2 Payload 序列化/反序列化

用于从 DomainEvent 中提取并转换 Payload 为具体的业务结构体。

使用场景:

  • Query Side: 处理具体的业务载荷

方法列表:

  • UnmarshalPayload[T any](ev BaseEvent) (T, error) - 从 DomainEvent 提取并反序列化 Payload
  • MarshalPayload(payload interface{}) ([]byte, error) - 序列化 Payload(特殊场景)

优势:

  • ✅ 统一使用 jxtjson(基于 jsoniter,性能比标准库快 2-3 倍)
  • ✅ 类型安全(泛型)
  • ✅ 自动处理 interface{} → map[string]interface{} → 结构体 的转换
  • ✅ 支持 []byte 和 RawMessage 类型的 Payload
  • ✅ 统一的错误处理
  • ✅ 避免各服务重复实现
  • ✅ 与 encoding/json 完全兼容

使用示例:

// 定义 Payload 结构
type ArchiveCreatedPayload struct {
    Title     string    `json:"title"`
    CreatedBy string    `json:"createdBy"`
    CreatedAt time.Time `json:"createdAt"`
}

// Query Side: 从 DomainEvent 提取 Payload
payload, err := jxtevent.UnmarshalPayload[ArchiveCreatedPayload](domainEvent)
if err != nil {
    return fmt.Errorf("failed to unmarshal payload: %w", err)
}

// 使用 payload
fmt.Println(payload.Title)

重要说明:

当 DomainEvent 从 JSON 反序列化时,Payload 字段的实际类型是 map[string]interface{}(而不是原始的结构体类型)。UnmarshalPayload 会自动处理这种情况,将 map 转换为目标结构体。

// 完整的端到端流程示例
// 1. Command Side: 创建事件
originalPayload := ArchiveCreatedPayload{Title: "Test", CreatedBy: "user-001"}
event := jxtevent.NewEnterpriseDomainEvent("Archive.Created", "archive-123", "Archive", originalPayload)

// 2. Command Side: 序列化保存到 Outbox
eventBytes, _ := jxtevent.MarshalDomainEvent(event)

// 3. Query Side: 从消息队列接收并反序列化
receivedEvent, _ := jxtevent.UnmarshalDomainEvent[*jxtevent.EnterpriseDomainEvent](eventBytes)
// 此时 receivedEvent.Payload 的类型是 map[string]interface{}

// 4. Query Side: 提取 Payload
payload, _ := jxtevent.UnmarshalPayload[ArchiveCreatedPayload](receivedEvent)
// 现在 payload 是 ArchiveCreatedPayload 结构体
4. ValidateConsistency(一致性校验)

校验 Envelope 与 DomainEvent 的一致性,确保事件在传输过程中关键信息不被篡改或不一致。

校验项:

  1. EventType 一致性
  2. AggregateID 一致性
  3. TenantId 一致性(如果是企业级事件)

使用示例:

// 创建 Envelope
envelope := &jxtevent.Envelope{
    EventType:   event.GetEventType(),
    AggregateID: event.GetAggregateID(),
    TenantID:    event.GetTenantId(),
    Payload:     payloadBytes,
}

// 校验一致性
if err := jxtevent.ValidateConsistency(envelope, event); err != nil {
    return fmt.Errorf("consistency validation failed: %w", err)
}

接口定义

BaseEvent 接口

所有事件都必须实现此接口:

type BaseEvent interface {
    GetEventID() string
    GetEventType() string
    GetOccurredAt() time.Time
    GetVersion() int
    GetAggregateID() string
    GetAggregateType() string
    GetPayload() interface{}
}
EnterpriseEvent 接口

多租户系统和需要可观测性支持的系统应实现此接口:

type EnterpriseEvent interface {
    BaseEvent
    
    // 租户隔离
    GetTenantId() int
    SetTenantId(int)
    
    // 可观测性方法
    GetCorrelationId() string
    SetCorrelationId(string)
    GetCausationId() string
    SetCausationId(string)
    GetTraceId() string
    SetTraceId(string)
}

最佳实践

1. 直接导入,避免间接层

✅ 推荐:

import jxtevent "github.com/ChenBigdata421/jxt-core/sdk/pkg/domain/event"

func Handle(evt *jxtevent.EnterpriseDomainEvent) error {
    // 直接使用,清晰明了
}

❌ 不推荐:

// 不要创建类型别名,增加间接层
type DomainEvent = jxtevent.EnterpriseDomainEvent
2. 使用泛型助手简化代码

✅ 推荐:

payload, err := jxtevent.UnmarshalPayload[MediaUploadedPayload](domainEvent)
if err != nil {
    return fmt.Errorf("failed to unmarshal payload: %w", err)
}

❌ 不推荐:

var json = jsoniter.ConfigCompatibleWithStandardLibrary
payloadBytes, _ := json.Marshal(domainEvent.GetPayload())
var payload MediaUploadedPayload
json.Unmarshal(payloadBytes, &payload)
3. 充分利用可观测性字段
event := jxtevent.NewEnterpriseDomainEvent(
    "Archive.Created",
    archiveID,
    "Archive",
    payload,
)

// 设置租户ID
event.SetTenantId(ctx.TenantID)

// 设置可观测性字段
event.SetCorrelationId(ctx.CorrelationID)  // 业务流程追踪
event.SetCausationId(triggerEventID)       // 因果事件ID
event.SetTraceId(ctx.TraceID)              // 分布式追踪ID
4. 事件命名规范

✅ 推荐(使用过去时态):

  • Archive.Created
  • Media.Uploaded
  • Relation.Deleted

❌ 不推荐(命令式):

  • CreateArchive
  • UploadMedia
5. Payload 设计原则

✅ 好的 Payload 设计:

type ArchiveCreatedPayload struct {
    ArchiveID   string    `json:"archiveId"`
    Title       string    `json:"title"`
    CreatedBy   string    `json:"createdBy"`
    CreatedAt   time.Time `json:"createdAt"`
}

❌ 避免包含过多信息:

type ArchiveCreatedPayload struct {
    Archive     *Archive  // 整个聚合根对象,太重
    User        *User     // 关联对象,不必要
    Permissions []string  // 权限信息,不属于事件
}

性能特性

JSON 序列化性能

本包使用统一的 jxtjson 包(基于 jsoniter),提供高性能的 JSON 序列化:

操作 性能指标 说明
序列化 ~690ns/op 比 encoding/json 快 2-3 倍
反序列化 ~1.2μs/op 比 encoding/json 快 2-3 倍
大 Payload ~511μs (1000 字段) 51KB JSON 数据
并发安全 ✅ 100 goroutines 无竞态条件
性能测试
# 运行性能基准测试
cd jxt-core/tests/domain/event/function_regression_tests
go test -run TestEnterpriseDomainEvent_PerformanceBenchmark -v

# 运行并发测试
go test -run TestEnterpriseDomainEvent_ConcurrentSerialization -v

测试

运行单元测试
cd jxt-core/sdk/pkg/domain/event
go test -v
运行回归测试
cd jxt-core/tests/domain/event/function_regression_tests
go test -v
测试覆盖率
  • 基础功能测试: 14 个测试用例
  • 企业级事件测试: 15 个测试用例
  • 序列化测试: 21 个测试用例
  • 集成测试: 9 个测试用例
  • Payload 测试: 13 个测试用例
  • 验证测试: 16 个测试用例

详细测试报告:

版本历史

  • v1.1.0 (2025-10-26): 序列化增强版本

    • ✅ 新增 21 个 EnterpriseDomainEvent 序列化/反序列化测试
    • ✅ 统一使用 jxtjson 包(基于 jsoniter)
    • ✅ 性能优化:序列化 ~690ns/op,反序列化 ~1.2μs/op
    • ✅ 完整的性能基准测试和并发测试
    • ✅ 与 encoding/json 完全兼容
    • ✅ 支持特殊字符(中文、Emoji、转义字符)
    • ✅ 支持大 Payload(1000+ 字段)
    • ✅ 完整的错误处理测试
  • v1.0.0 (2025-10-25): 初始版本

    • 实现 BaseDomainEvent
    • 实现 EnterpriseDomainEvent
    • 实现 UnmarshalPayload 泛型助手
    • 实现 ValidateConsistency 一致性校验
    • 完整的单元测试覆盖

相关文档

核心文档
测试文档
迁移文档

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func MarshalDomainEvent

func MarshalDomainEvent(event BaseEvent) ([]byte, error)

MarshalDomainEvent 序列化完整的 DomainEvent 用于 Command Side 保存到 Outbox 或发布到消息队列

使用示例:

eventBytes, err := event.MarshalDomainEvent(domainEvent)
if err != nil {
    return fmt.Errorf("failed to marshal domain event: %w", err)
}

优势: 1. 统一使用 jsoniter.ConfigCompatibleWithStandardLibrary 2. 统一的错误处理 3. 避免各服务重复实现

注意: - 序列化时会包含所有字段(EventID, EventType, Payload 等) - Payload 字段会被序列化为 JSON 对象(嵌套在 DomainEvent JSON 中)

func MarshalDomainEventToString

func MarshalDomainEventToString(event BaseEvent) (string, error)

MarshalDomainEventToString 序列化 DomainEvent 为 JSON 字符串 用于测试或特殊场景

使用示例:

jsonString, err := event.MarshalDomainEventToString(domainEvent)
if err != nil {
    return fmt.Errorf("failed to marshal domain event: %w", err)
}

func MarshalPayload

func MarshalPayload(payload interface{}) ([]byte, error)

MarshalPayload 标准化的Payload序列化助手(用于特殊场景) 一般不需要直接使用,UnmarshalPayload 内部会自动处理

使用示例:

payloadBytes, err := event.MarshalPayload(payload)
if err != nil {
    return fmt.Errorf("failed to marshal payload: %w", err)
}

优势: 1. 统一使用 jxt-core/sdk/pkg/json 的 jsoniter 配置 2. 统一的错误处理

func UnmarshalDomainEvent

func UnmarshalDomainEvent[T BaseEvent](data []byte) (T, error)

UnmarshalDomainEvent 反序列化完整的 DomainEvent 用于 Query Side 从消息队列接收事件

使用示例:

// 使用 BaseDomainEvent
domainEvent, err := event.UnmarshalDomainEvent[*event.BaseDomainEvent](msg.Payload)
if err != nil {
    return fmt.Errorf("failed to unmarshal domain event: %w", err)
}

// 使用 EnterpriseDomainEvent
enterpriseEvent, err := event.UnmarshalDomainEvent[*event.EnterpriseDomainEvent](msg.Payload)
if err != nil {
    return fmt.Errorf("failed to unmarshal enterprise event: %w", err)
}

优势: 1. 统一使用 jsoniter.ConfigCompatibleWithStandardLibrary 2. 类型安全(泛型) 3. 统一的错误处理 4. 避免各服务重复实现

注意: - 反序列化后,Payload 字段的类型是 map[string]interface{}(而不是原始结构体) - 需要使用 UnmarshalPayload 进一步提取具体的 Payload 结构体

func UnmarshalDomainEventFromString

func UnmarshalDomainEventFromString[T BaseEvent](jsonString string) (T, error)

UnmarshalDomainEventFromString 从 JSON 字符串反序列化 DomainEvent 用于测试或特殊场景

使用示例:

domainEvent, err := event.UnmarshalDomainEventFromString[*event.BaseDomainEvent](jsonString)
if err != nil {
    return fmt.Errorf("failed to unmarshal domain event: %w", err)
}

func UnmarshalPayload

func UnmarshalPayload[T any](ev BaseEvent) (T, error)

UnmarshalPayload 标准化的Payload反序列化助手 解决各服务自行处理序列化导致的不一致问题

使用示例:

payload, err := event.UnmarshalPayload[MediaUploadedPayload](domainEvent)
if err != nil {
    return fmt.Errorf("failed to unmarshal payload: %w", err)
}

优势: 1. 统一使用jsoniter.ConfigCompatibleWithStandardLibrary 2. 类型安全(泛型) 3. 自动处理 interface{} → map[string]interface{} → 结构体 的转换 4. 支持 []byte 和 RawMessage 类型的 Payload 5. 统一的错误处理 6. 避免各服务重复实现

注意: - 当 Payload 是 interface{} 类型且从 JSON 反序列化时,实际类型是 map[string]interface{} - 本方法会自动处理这种情况,将 map 转换为目标结构体

func ValidateConsistency

func ValidateConsistency(envelope *Envelope, event BaseEvent) error

ValidateConsistency 校验Envelope与DomainEvent的一致性 确保事件在传输过程中关键信息不被篡改或不一致

校验项: 1. EventType一致性 2. AggregateID一致性 3. TenantId一致性(如果是企业级事件)

使用场景: - Outbox适配器在保存事件前校验 - EventHandler在处理事件前校验 - 测试用例中验证事件完整性

Types

type BaseDomainEvent

type BaseDomainEvent struct {
	// ========== 技术基础字段 ==========
	EventID    string    `json:"eventId" gorm:"type:char(36);primary_key;column:event_id;comment:事件ID"`
	EventType  string    `json:"eventType" gorm:"type:varchar(255);index;column:event_type;comment:事件类型"`
	OccurredAt time.Time `json:"occurredAt" gorm:"type:datetime;index;column:occurred_at;comment:事件发生时间"`
	Version    int       `json:"version" gorm:"type:int;column:event_version;comment:事件版本"`

	// ========== DDD核心概念 ==========
	AggregateID   string `json:"aggregateId" gorm:"type:varchar(255);index;column:aggregate_id;comment:聚合根ID"`
	AggregateType string `json:"aggregateType" gorm:"type:varchar(255);index;column:aggregate_type;comment:聚合根类型"`

	// ========== 事件载荷 ==========
	Payload interface{} `json:"payload" gorm:"-"` // 不持久化到数据库
}

BaseDomainEvent 基础领域事件结构 包含所有事件驱动系统都需要的核心字段 适用于:所有使用事件驱动架构的系统

func NewBaseDomainEvent

func NewBaseDomainEvent(eventType string, aggregateID interface{}, aggregateType string, payload interface{}) *BaseDomainEvent

NewBaseDomainEvent 创建基础领域事件

func (*BaseDomainEvent) GetAggregateID

func (e *BaseDomainEvent) GetAggregateID() string

func (*BaseDomainEvent) GetAggregateType

func (e *BaseDomainEvent) GetAggregateType() string

func (*BaseDomainEvent) GetEventID

func (e *BaseDomainEvent) GetEventID() string

Getter方法

func (*BaseDomainEvent) GetEventType

func (e *BaseDomainEvent) GetEventType() string

func (*BaseDomainEvent) GetOccurredAt

func (e *BaseDomainEvent) GetOccurredAt() time.Time

func (*BaseDomainEvent) GetPayload

func (e *BaseDomainEvent) GetPayload() interface{}

func (*BaseDomainEvent) GetVersion

func (e *BaseDomainEvent) GetVersion() int

type BaseEvent

type BaseEvent interface {
	GetEventID() string
	GetEventType() string
	GetOccurredAt() time.Time
	GetVersion() int
	GetAggregateID() string
	GetAggregateType() string
	GetPayload() interface{}
}

BaseEvent 基础事件接口 所有事件都必须实现此接口

type EnterpriseDomainEvent

type EnterpriseDomainEvent struct {
	BaseDomainEvent

	// ========== 企业级通用字段 ==========
	// 租户隔离:多租户系统的核心字段
	// 多租户系统:使用实际的租户ID(如 1, 2, 3)
	// 单租户系统:使用 0 表示全局/无租户
	TenantId int `json:"tenantId" gorm:"type:int;index;column:tenant_id;comment:租户ID"`

	// ========== 可观测性字段 ==========
	// 用于分布式追踪和因果链路分析
	CorrelationId string `json:"correlationId,omitempty" gorm:"type:varchar(255);index;column:correlation_id;comment:业务关联ID"`
	CausationId   string `json:"causationId,omitempty" gorm:"type:varchar(255);index;column:causation_id;comment:因果事件ID"`
	TraceId       string `json:"traceId,omitempty" gorm:"type:varchar(255);index;column:trace_id;comment:分布式追踪ID"`
}

EnterpriseDomainEvent 企业级领域事件结构 在BaseDomainEvent基础上增加企业级通用字段 适用于:多租户SaaS系统、企业级应用

func NewEnterpriseDomainEvent

func NewEnterpriseDomainEvent(eventType string, aggregateID interface{}, aggregateType string, payload interface{}) *EnterpriseDomainEvent

NewEnterpriseDomainEvent 创建企业级领域事件

func (*EnterpriseDomainEvent) GetCausationId

func (e *EnterpriseDomainEvent) GetCausationId() string

func (*EnterpriseDomainEvent) GetCorrelationId

func (e *EnterpriseDomainEvent) GetCorrelationId() string

可观测性方法

func (*EnterpriseDomainEvent) GetTenantId

func (e *EnterpriseDomainEvent) GetTenantId() int

企业级特定方法

func (*EnterpriseDomainEvent) GetTraceId

func (e *EnterpriseDomainEvent) GetTraceId() string

func (*EnterpriseDomainEvent) SetCausationId

func (e *EnterpriseDomainEvent) SetCausationId(id string)

func (*EnterpriseDomainEvent) SetCorrelationId

func (e *EnterpriseDomainEvent) SetCorrelationId(id string)

func (*EnterpriseDomainEvent) SetTenantId

func (e *EnterpriseDomainEvent) SetTenantId(id int)

func (*EnterpriseDomainEvent) SetTraceId

func (e *EnterpriseDomainEvent) SetTraceId(id string)

type EnterpriseEvent

type EnterpriseEvent interface {
	BaseEvent

	// 租户隔离
	GetTenantId() int
	SetTenantId(int)

	// 可观测性方法
	GetCorrelationId() string
	SetCorrelationId(string)
	GetCausationId() string
	SetCausationId(string)
	GetTraceId() string
	SetTraceId(string)
}

EnterpriseEvent 企业级事件接口 多租户系统和需要可观测性支持的系统应实现此接口

type Envelope

type Envelope struct {
	EventType   string
	AggregateID string
	TenantID    int
	Payload     []byte
}

Envelope 事件信封结构 用于事件传输和一致性校验

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL