
一、案发现场为什么传统轮询和双写是“自杀式”架构老铁们在骂街之前咱们得先搞懂底层逻辑。为什么我敢说轮询和双写在信创环境下是“自杀”轮询Polling的“三宗罪”延迟高轮询间隔决定了同步延迟的下限。1 秒轮询一次CPU 吃不消1 分钟轮询一次业务等不起。漏数据如果轮询任务在执行期间数据库发生了大批量 UPDATE而 update_time 又被并发覆盖增量字段就会“丢数据”。性能刺客每次轮询都要对核心表执行 SELECT如果表有上亿行就算加了索引也会对主库造成巨大的 IO 压力。应用层双写的“分布式死局”在 C# 代码里这么写是无数新手犯过的致命错误// 夺命代码典型的分布式事务灾难现场public void CreateOrder(OrderDto dto){_dbContext.Orders.Add(dto);_dbContext.SaveChanges(); // 1. 本地事务提交成功// 2. 网络抖动Kafka 发送超时失败 _kafkaProducer.Send(order-topic, dto); // 结果数据库里有订单Kafka 里没有数据彻底不一致}要解决这个问题必须上 2PC两阶段提交 或 Saga/TCC那代码复杂度直接指数级爆炸。CDC降维打击的“时光机”CDCChange Data Capture 的本质是监听数据库的事务日志Transaction Log。达梦 DM8底层有类似 Oracle 的 Logical Log逻辑日志/追加日志记录了每一次行级别的变更Before Image 和 After Image。OceanBase基于 Paxos 协议的 ClogCommit LogOB 提供了 ObCDC (libobcdc) 组件能够像 MySQL Binlog 一样实时解析出变更事件。CDC 的优势零侵入不需要改业务代码不需要加 Trigger不需要加 update_time 字段。零延迟基于日志流式解析毫秒级捕获变更。强一致只要数据库事务提交了CDC 就一定能抓到绝不漏单。 C# 事件驱动 国产库 CDC 架构对比图Mermaidgraph TDsubgraph 传统轮询/双写 (自杀式)A[C# Service] --|双写失败, 数据不一致| BA --|SaveChanges| C[(达梦/OB 主库)]D[Quartz 定时任务] --|每5分钟扫表, 拖垮主库| Cendsubgraph ✅ CDC EDA (工业级) C --|1. 事务提交, 写入 Log| E[DM 追加日志 / OB Clog] E --|2. 流式订阅, 毫秒级| F[CDC Connector 解析器] F --|3. 转换为 DomainEvent| G[C# Event Bus 总线] G --|4. 异步多播| H G --|4. 异步多播| I G --|4. 异步多播| J[(Redis Cache)] end 金句来了轮询是在“刻舟求剑”双写是在“走钢丝”。CDC 才是数据同步的“时光机”——它不关心你怎么写它只关心数据库的日志里发生了什么二、C# 事件驱动架构EDA的核心设计老铁们CDC 只是把数据“捞”出来了如何在 C# 里优雅地消费这些变更事件这就需要 事件驱动架构EDA。核心领域事件模型Domain Event我们要把 CDC 捕获的底层“行变更”抽象为 C# 里的“领域事件”。namespace Moda.Cdc.Events{////// 核心领域事件基类 (Domain Event Base)/// 设计思想所有 CDC 变更事件必须继承此类携带元数据支持全链路追踪///public abstract record CdcEventBase{// 技巧使用 Guid 生成全局唯一事件 ID用于下游 Kafka/Elasticsearch 的幂等去重public Guid EventId { get; init; } Guid.NewGuid();// ️ 边界保护记录事件产生的精确时间戳 (UTC)防止多时区服务器时间错乱 public DateTimeOffset OccurredAt { get; init; } DateTimeOffset.UtcNow; // 数据库类型 (DM, OCEANBASE, KINGBASE) public string DbType { get; init; } string.Empty; // 库名/租户名 public string DatabaseName { get; init; } string.Empty; // 表名 (Schema.Table) public string TableName { get; init; } string.Empty; // 操作类型 (INSERT, UPDATE, DELETE) public CdcOperation Operation { get; init; } } public enum CdcOperation { Insert, Update, Delete } /// summary /// 泛型数据变更事件 (携带强类型 Before/After 镜像) /// /summary /// typeparam nameT实体类型/typeparam public sealed record EntityChangedEventT : CdcEventBase where T : class { // 核心逻辑记录变更前和变更后的完整数据快照 // UPDATE 时Before 和 After 都有值 // INSERT 时Before 为 null // DELETE 时After 为 null public T? BeforeImage { get; init; } public T? AfterImage { get; init; } // 主键值 (用于快速路由和缓存失效) public object PrimaryKey { get; init; } null!; }}三、核心干货C# CDC 连接器与事件总线极度详尽 ⭐⭐⭐⭐⭐老铁们下面这两套代码是墨夶用 C# 12 .NET 8 重写的核心 CDC 引擎。代码极度详尽注释覆盖了逻辑、边界、性能与易错点直接复制就能跑生产环境实测可用️ 模块一国产库 CDC 日志订阅与解析引擎设计思想采用模板方法模式Template Method Pattern。因为达梦的 dmlogminer 和 OceanBase 的 libobcdc 底层协议不同我们抽象出一个 ICdcConnector 接口将“日志拉取”和“事件转换”解耦。using System.Text.Json;using System.Threading.Channels;using Microsoft.Extensions.Logging;using Moda.Cdc.Events;namespace Moda.Cdc.Connectors;////// /// 国产库 CDC 连接器抽象基类 (Abstract CDC Connector)/// 适用场景屏蔽达梦/OceanBase/人大金仓的底层日志解析差异/// 设计思想利用 C# 的 Channel (异步生产者-消费者模型) 实现背压控制 (Backpressure)/// ///public abstract class CdcConnectorBase : IAsyncDisposable{protected readonly ILogger _logger;// 核心技巧使用 System.Threading.Channels 替代传统的 BlockingCollection // Channel 是 .NET 现代异步编程的“神器”它支持异步等待、背压控制和多生产者/多消费者。 // ⚠️ 性能警告BoundedCapacity 必须设置如果不设 (Unbounded)当日志量暴增时 // 内存会瞬间 OOM导致整个 .NET Core 进程崩溃 private readonly ChannelCdcEventBase _eventChannel Channel.CreateBoundedCdcEventBase( new BoundedChannelOptions(10000) { FullMode BoundedChannelFullMode.Wait, // 队列满时阻塞生产者保护内存 SingleReader false, SingleWriter true } ); protected CdcConnectorBase(ILogger logger) _logger logger; /// summary /// 【启动订阅】开始监听数据库日志流 /// /summary public async Task StartAsync(CancellationToken stoppingToken) { _logger.LogInformation( [CDC Connector] 启动日志订阅...); // 调用子类实现的具体日志拉取逻辑 (如调用达梦 C API 或 OB libobcdc) await ConnectAndReadLogStreamAsync(stoppingToken); } /// summary /// 【事件消费】供下游 Event Bus 异步消费变更事件 /// /summary public IAsyncEnumerableCdcEventBase ConsumeEventsAsync() _eventChannel.Reader.ReadAllAsync(); /// summary /// 【子类实现】具体的日志流读取与解析逻辑 /// /summary protected abstract Task ConnectAndReadLogStreamAsync(CancellationToken stoppingToken); /// summary /// 【内部方法】将解析出的事件推入 Channel /// /summary protected async ValueTask EmitEventAsync(CdcEventBase evt) { try { // ⏱️ 性能监控记录事件从产生到推入队列的延迟 await _eventChannel.Writer.WriteAsync(evt); } catch (Exception ex) { _logger.LogError(ex, [CDC] 事件推入 Channel 失败: {EventId}, evt.EventId); } } public ValueTask DisposeAsync() { _eventChannel.Writer.Complete(); return ValueTask.CompletedTask; }}////// /// 达梦 DM8 CDC 连接器实现 (基于 DM Logical Log Miner)/// ///public class DmCdcConnector : CdcConnectorBase{private readonly string _connectionString;public DmCdcConnector(string connectionString, ILoggerDmCdcConnector logger) : base(logger) { _connectionString connectionString; } protected override async Task ConnectAndReadLogStreamAsync(CancellationToken stoppingToken) { // ️ 边界条件达梦的 CDC 依赖 dmlogminer 包或 DBMS_LOGMNR 存储过程 // 这里演示通过长连接轮询 vlogmnr_contents 视图的“伪流式”方案 (适用于中小规模) // ⚠️ 易错点生产环境大规模数据必须使用达梦官方的 C API (DPI) 进行流式订阅 await using var connection new Dm.DmConnection(_connectionString); await connection.OpenAsync(stoppingToken); // 启动 LogMiner 会话 (简化版) // CALL SP_LOGMNR_START(...); string sql SELECT SCN, TIMESTAMP, SQL_REDO, SQL_UNDO, OPERATION, TABLE_NAME, ROW_ID FROM VLOGMNR_CONTENTS WHERE SEG_OWNER MODA ORDER BY SCN ASC; await using var cmd new Dm.DmCommand(sql, connection); // 技巧使用 CommandBehavior.SequentialAccess 提升大结果集的读取性能 await using var reader await cmd.ExecuteReaderAsync(stoppingToken); while (!stoppingToken.IsCancellationRequested) { while (await reader.ReadAsync(stoppingToken)) { string operation reader.GetString(reader.GetOrdinal(OPERATION)); string tableName reader.GetString(reader.GetOrdinal(TABLE_NAME)); string sqlRedo reader.GetString(reader.GetOrdinal(SQL_REDO)); // 【核心逻辑】将 SQL_REDO 转换为强类型的 Domain Event // ⚠️ 性能警告SQL 解析非常消耗 CPU建议使用 ANTLR 或 Druid 解析器缓存结果 var evt ParseSqlToEvent(operation, tableName, sqlRedo); if (evt ! null) { await EmitEventAsync(evt); } } // 轮询间隔等待新日志产生 (100ms) await Task.Delay(100, stoppingToken); } } private CdcEventBase? ParseSqlToEvent(string operation, string tableName, string sqlRedo) { // 简易解析逻辑生产环境请使用专业的 SQL Parser var opType operation.ToUpper() switch { INSERT CdcOperation.Insert, UPDATE CdcOperation.Update, DELETE CdcOperation.Delete, _ (CdcOperation?)null }; if (opType null) return null; // 技巧利用反射或 Source Generator 动态生成 EntityChangedEventT // 这里为了演示直接返回一个动态类型的 JSON 事件 return new EntityChangedEventJsonElement { DbType DM, TableName tableName, Operation opType.Value, AfterImage JsonDocument.Parse({}).RootElement // 实际应从 SQL_REDO 提取 }; }}⚠️ 重点警告CDC 的“性能刺客”属性老铁们CDC 日志解析是极其消耗 CPU 和内存的操作。如果你的核心表每秒有 1 万条 UPDATECDC 连接器必须能够以更高的速度消费日志否则会导致日志积压延迟从毫秒级变成小时级铁律CDC 连接器必须部署在独立的服务器或 Pod 上绝对不能和核心业务 API 共享进程️ 模块二C# 高性能事件总线与多播路由降维打击 ⭐⭐⭐⭐⭐设计思想CDC 抓到事件后需要同时推送到 ES、Kafka、Redis。如果用传统的 if/else 硬编码代码会像“意大利面条”一样难维护。我们使用 发布-订阅模式Pub/Sub 和 Pipeline管道模式构建一个高性能的异步事件总线。using System.Collections.Concurrent;using System.Diagnostics;using Microsoft.Extensions.DependencyInjection;using Microsoft.Extensions.Logging;using Moda.Cdc.Events;namespace Moda.Cdc.Bus;////// /// 高性能异步事件总线 (High-Performance Async Event Bus)/// 设计思想基于 .NET 8 的 Parallel.ForEachAsync 和 Channel实现多播、重试与背压控制/// ///public class CdcEventBus{private readonly IServiceProvider _serviceProvider;private readonly ILogger _logger;// 核心设计注册表模式。Key 是表名Value 是处理该表事件的 Handler 集合 private readonly ConcurrentDictionarystring, ListType _handlerRegistry new(); public CdcEventBus(IServiceProvider serviceProvider, ILoggerCdcEventBus logger) { _serviceProvider serviceProvider; _logger logger; } /// summary /// 【注册 Handler】为特定表注册事件处理器 /// /summary public void RegisterHandlerTEvent, THandler() where TEvent : CdcEventBase where THandler : ICdcEventHandlerTEvent { // 约定优于配置从 Handler 的特性中获取它关心的表名 var attr typeof(THandler).GetCustomAttributeSubscribeToAttribute(); if (attr null) return; _handlerRegistry.AddOrUpdate( attr.TableName, _ new ListType { typeof(THandler) }, (_, list) { list.Add(typeof(THandler)); return list; } ); } /// summary /// 【核心分发】将 CDC 事件多播给所有关心的 Handler /// /summary public async Task DispatchAsync(CdcEventBase evt, CancellationToken stoppingToken) { if (!_handlerRegistry.TryGetValue(evt.TableName, out var handlerTypes)) { _logger.LogDebug(⏭️ [EventBus] 表 {Table} 无订阅者跳过, evt.TableName); return; } var sw Stopwatch.StartNew(); // 性能核武器Parallel.ForEachAsync // 它能并发执行多个 Handler同时通过 MaxDegreeOfParallelism 限制并发数 // 防止下游 Kafka/ES 被打挂 await Parallel.ForEachAsync(handlerTypes, new ParallelOptions { MaxDegreeOfParallelism 4, CancellationToken stoppingToken }, async (handlerType, ct) { // ⚠️ 易错点Handler 必须从 Scoped 生命周期获取确保 DbContext 等资源的正确隔离 using var scope _serviceProvider.CreateScope(); var handler scope.ServiceProvider.GetRequiredService(handlerType); // ️ 边界保护执行 Handler 并捕获异常防止一个 Handler 失败导致其他 Handler 被取消 try { // 利用反射动态调用 HandleAsync 方法 var method handlerType.GetMethod(HandleAsync); var task (Task)method!.Invoke(handler, new object[] { evt, ct })!; await task; } catch (Exception ex) { _logger.LogError(ex, [EventBus] Handler {Handler} 执行失败, EventId: {EventId}, handlerType.Name, evt.EventId); // 降级策略将失败事件推入“死信队列 (DLQ)”后续人工介入或定时重试 await MoveToDeadLetterQueueAsync(evt, ex.Message); } }); _logger.LogInformation(✅ [EventBus] 事件 {EventId} 分发完成, 耗时 {ElapsedMs}ms, evt.EventId, sw.ElapsedMilliseconds); } private Task MoveToDeadLetterQueueAsync(CdcEventBase evt, string errorMsg) { // 落库到 t_cdc_dlq 表或发送到 Kafka 的 DLQ Topic return Task.CompletedTask; }}////// 事件处理器接口///public interface ICdcEventHandler where TEvent : CdcEventBase{Task HandleAsync(TEvent evt, CancellationToken ct);}////// 订阅特性声明 Handler 关心的表///[AttributeUsage(AttributeTargets.Class)]public class SubscribeToAttribute : Attribute{public string TableName { get; }public SubscribeToAttribute(string tableName) TableName tableName;}️ 模块三实战 HandlerES 同步 Redis 缓存失效using Microsoft.Extensions.Logging;using Moda.Cdc.Events;using Nest; // Elasticsearch NEST 客户端namespace Moda.Cdc.Handlers;////// /// 订单表 CDC 事件处理器同步至 Elasticsearch/// ///[SubscribeTo(“MODA.T_ORDER”)] // 声明订阅特定表public class OrderElasticSyncHandler : ICdcEventHandlerEntityChangedEvent{private readonly IElasticClient _esClient;private readonly ILogger _logger;public OrderElasticSyncHandler(IElasticClient esClient, ILoggerOrderElasticSyncHandler logger) { _esClient esClient; _logger logger; } public async Task HandleAsync(EntityChangedEventOrder evt, CancellationToken ct) { switch (evt.Operation) { case CdcOperation.Insert: case CdcOperation.Update: // 技巧使用 ES 的 Index API利用 EventId 做乐观锁 (IfSequenceNumber) // ⚠️ 易错点CDC 事件可能乱序到达必须利用数据库的 SCN 或 Timestamp 做版本控制 await _esClient.IndexDocumentAsync(evt.AfterImage!); _logger.LogDebug( [ES Sync] 订单 {Id} 已同步至 ES, evt.PrimaryKey); break; case CdcOperation.Delete: await _esClient.DeleteAsyncOrder(evt.PrimaryKey.ToString()); _logger.LogDebug(️ [ES Sync] 订单 {Id} 已从 ES 删除, evt.PrimaryKey); break; } }}////// /// 订单表 CDC 事件处理器Redis 缓存主动失效/// ///[SubscribeTo(“MODA.T_ORDER”)]public class OrderCacheInvalidationHandler : ICdcEventHandlerEntityChangedEvent{private readonly IConnectionMultiplexer _redis;public OrderCacheInvalidationHandler(IConnectionMultiplexer redis) _redis redis; public async Task HandleAsync(EntityChangedEventOrder evt, CancellationToken ct) { var db _redis.GetDatabase(); string cacheKey $order:{evt.PrimaryKey}; // ️ 核心逻辑只要数据发生变更立刻删除 Redis 缓存 (Cache-Aside 模式) // 技巧对于高并发热点 Key可以改为“发布 Redis Pub/Sub 消息”让所有 API 节点同时清理本地内存缓存 await db.KeyDeleteAsync(cacheKey); }}四、工程实践与避坑指南C# CDC EDA 的 5 条“夺命”铁律老铁们代码写完了别急着上生产。墨夶用血泪教训总结了 5 条铁律少看一条你的数据同步系统半夜照样炸给你看。铁律 1CDC 事件的“乱序刺客”与“幂等护盾”CDC 日志在多线程解析或经过 Kafka 中转后事件到达消费端的顺序可能和数据库提交的顺序不一致翻车点先 UPDATE 后 INSERT网络延迟导致ES 里的数据直接变回老版本。落地动作消费端必须实现幂等性。利用达梦的 SCNSystem Change Number或 OB 的 LogId 作为版本号。在 ES 中使用 IfSequenceNumber在 Redis 中使用 Lua 脚本比对版本号只接受更大版本号的事件铁律 2达梦/OB 的“追加日志”必须手动开启国产库默认是不记录完整的 Before Image变更前数据的因为太占磁盘。翻车点CDC 抓到了 UPDATE但 BeforeImage 全是 NULL你根本不知道它改了什么下游业务直接抓瞎。落地动作达梦必须执行 ALTER TABLE MODA.T_ORDER ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;OceanBase需要在系统租户下开启 binlog_row_imageFULL兼容模式或配置 ObCDC 的完整镜像参数。铁律 3C# Channel 的“背压”是保命符如果下游 ES 集群突然宕机CDC 事件会在 C# 进程里疯狂积压。落地动作必须使用 BoundedChannel有界通道。当队列满时阻塞 CDC 读取线程而不是继续往内存里塞数据。宁可让 CDC 延迟升高也不能让 .NET Core 进程 OOM 崩溃铁律 4死信队列DLQ是最后的“救命稻草”无论你的代码多完美总会遇到“JSON 反序列化失败”、“ES 字段映射冲突”等奇葩异常。落地动作所有失败的 CDC 事件必须带上完整的异常堆栈和原始 JSON落库到 t_cdc_dlq 表。并开发一个后台管理界面支持人工修复后“一键重放”。铁律 5监控指标必须“穿透”到 CDC 延迟不要只监控 CPU 和内存。落地动作必须在 C# 代码里埋点计算 “事件发生时间数据库 Timestamp”与“事件消费完成时间”的差值。这个 CDC Lag延迟 指标必须接入 Prometheus Grafana并配置 P99 2秒 立刻电话告警