Skip to content

Repository files navigation

项目概览

  • 名称: Xmart.CDC.MySql — 基于 MySQL binlog 的 CDC 读取库。
  • 用途: 从 MySQL 二进制日志(binlog)读取变更事件,支持独立的行级变更与 schema 变更订阅。
  • 核心模型: ChangeRecord 表示 INSERT/UPDATE/DELETE 行级变更,SchemaChangeRecord 表示 DDL/schema 变更。

主要功能

  • 解析 MySQL binlog row event 与 query event
  • 支持 IAsyncEnumerable<ChangeRecord> 流式消费
  • 支持 OnRowChange / OnSchemaChange 事件回调
  • 支持 OnCommitPosition 定期提交位置
  • 支持 TableMap 缓存失效,处理 DDL 后表结构变化

架构说明

  • BinlogReader 通过 TCP 连接读取 MySQL binlog,并按事件头解析每个 binlog 包
  • TableMapEvent 用于缓存表结构和列元数据,后续 row event 解析依赖它
  • QueryEvent 中的 DDL 语句会被解析为 SchemaChangeRecord
  • ChangeRecord 表示 INSERT / UPDATE / DELETE 行级变更,SchemaChangeRecord 表示 CREATE / ALTER / DROP / RENAME / TRUNCATE 等 schema 变更
  • OnRowChangeOnSchemaChange 为两条独立消费通道,避免把 schema 事件混入行级变更流
  • OnChange 仍保留为兼容旧 API 的过渡项,但新代码应优先使用 OnRowChange
  • DDL 事件后会触发 TableMapCache 失效,保证后续 row event 能重新加载最新表结构
  • 完整架构导览:请参阅 docs/architecture.md(运行时数据流、模块分层、关键机制、消费者 API 设计)
  • 未来演进路线:请参阅 docs/roadmap.md(M0 已完成 + M1-M4 的 33 个实现 ticket)

快速开始

示例代码

var reader = new BinlogReaderBuilder()
    .WithConnection("127.0.0.1", 3306, "root", "password")
    .WithDatabase("my_db")
    .WithTables("my_table")
    .WithServerId(1234)
    .FromCurrentPosition()
    .Build();

reader.OnRowChange += async record =>
{
    Console.WriteLine($"Row change: {record.EventType} {record.Database}.{record.Table} @{record.Position}");
    await Task.CompletedTask;
};

reader.OnSchemaChange += async schemaChange =>
{
    Console.WriteLine($"Schema change: {schemaChange.ChangeType} {schemaChange.Database}.{schemaChange.Table} query={schemaChange.Query}");
    await Task.CompletedTask;
};

await foreach (var record in reader.ReadAsync(cancellationToken))
{
    // 处理行级变更
}

运行示例

// Path: README.md
dotnet run --project samples/Xmart.CDC.MySql.Samples

请先修改 samples/Xmart.CDC.MySql.Samples/appsettings.json,让 MySql 连接信息指向你的 MySQL 实例。

运行测试

// Path: README.md
dotnet test --no-build

集成测试需要访问真实 MySQL 实例。可通过环境变量传入:

// Path: README.md
$env:XMART_CDC_TEST_CONNECTION = 'Server=127.0.0.1;Port=3306;Database=mysql;User=root;Password=...;'

编译与打包

// Path: README.md
dotnet build Xmart.CDC.slnx
dotnet pack src/Xmart.CDC.MySql -o nupkg

注意

  • 推荐使用 OnRowChange 处理行事件,OnSchemaChange 处理 DDL/schema 事件。
  • OnChange 仍作为兼容旧用户的过渡事件,建议迁移至 OnRowChange

数据类型支持矩阵 ChangeRecord.Before/AfterRowData 提供类型化值,覆盖常用 MySQL 数据类型:

类别 MySQL 类型 .NET 类型
整数 TINYINT / SMALLINT / MEDIUMINT / INT / BIGINT sbyte / short / int / long
无符号整数 * UNSIGNED byte / ushort / uint / ulong
浮点 FLOAT / DOUBLE float / double
定点数 DECIMAL / NUMERIC decimal
字符串 VARCHAR / CHAR / TEXT / ENUM(索引) / SET(位图) string / long / ulong
二进制 BLOB / BINARY byte[]
时间 DATE / TIME / DATETIME / TIMESTAMP / YEAR DateOnly / TimeOnly / DateTime / int
BIT(n) long

说明:ENUM 返回成员索引(1-based),SET 返回位图整数;列名与 UNSIGNED 标志从 information_schema 延迟解析(详见 ADR-0007)。

列名解析行为

  • 收到某表的首个 TABLE_MAP 事件时,自动查询 information_schema.COLUMNS 获取真实列名与 UNSIGNED 标志
  • 结果缓存;检测到 DDL(QueryEvent)时自动失效并重新查询
  • 查询失败时静默降级为 col_X 占位列名,不影响主流程

数据类型覆盖演示 运行示例菜单第 2 项(TypeCoverageDemo)可端到端验证 20+ 种数据类型解析正确性:自动建表 → 插入已知值 → CDC 捕获 → 逐列比对 → 清理临时表。

常见问题与排障

  • 无法捕获事件:请检查 MySQL binlog 是否开启、账号是否有 REPLICATION SLAVE / REPLICATION CLIENT 权限。
  • DDL 事件后若出现表结构错误:本库会尝试在 OnSchemaChange 后失效旧 TableMap 缓存,确保下一个 TABLE_MAP 事件重新加载。
  • ALTER ADD COLUMN 后新列为 NULL:正常行为——NULL 列保留在 RowData 中(值为 null),不缺失。

贡献指南

  • Fork → 新分支 → 提交(单一功能优先)→ 发起 PR。
  • 运行 dotnet test 并确保所有测试通过。

参考

如果你希望我把 README 翻译成英文版、添加更多使用示例或把常见错误的定位步骤写成检查清单,我可以继续完善。

About

No description, website, or topics provided.

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages