流式處理“的極度詳盡、滿載干貨的技術(shù))
一、為什么Kafka 國(guó)產(chǎn)庫(kù)總是一上線就翻車先搞懂死因才能對(duì)癥下藥。信創(chuàng)消息消費(fèi)翻車通常是消費(fèi)模式、寫入方式、事務(wù)控制三端集體拉胯。1.1 單條插入Single Insert是萬(wàn)惡之源很多老鐵寫消費(fèi)者習(xí)慣性地while(true) {var msg consumer.Consume();// 執(zhí)行 INSERT INTO …consumer.Commit(msg);}魔性比喻 這就像你要搬 10 萬(wàn)塊磚你每次只拿 1 塊還要跑 10 萬(wàn)趟每次 INSERT金倉(cāng)都要經(jīng)歷解析 SQL - 生成執(zhí)行計(jì)劃 - 開(kāi)啟隱式事務(wù) - 寫 WAL 日志 - 提交事務(wù) - 刷盤。10 萬(wàn) TPS金倉(cāng)每秒要開(kāi) 10 萬(wàn)個(gè)事務(wù)CPU 不爆才怪1.2 并發(fā)寫入的死鎖魔咒為了提速有人開(kāi)了 50 個(gè) Task 并發(fā)寫。坑點(diǎn) 人大金倉(cāng)基于 PostgreSQL 魔改的 MVCC 和鎖機(jī)制在多并發(fā) UPDATE 同一行比如累加設(shè)備在線狀態(tài)時(shí)極易產(chǎn)生死鎖。或者因?yàn)椴迦腠樞虿灰恢聦?dǎo)致 B 樹(shù)索引頁(yè)分裂性能斷崖式下跌。1.3 “至少一次”At-Least-Once的重復(fù)消費(fèi)陷阱網(wǎng)絡(luò)一抖動(dòng)Kafka 消費(fèi)者崩潰重啟后從上一個(gè) Offset 重新消費(fèi)。結(jié)果 同一批告警數(shù)據(jù)被寫進(jìn)金倉(cāng)兩次如果業(yè)務(wù)沒(méi)做冪等Idempotency 控制數(shù)據(jù)庫(kù)里全是重復(fù)數(shù)據(jù)報(bào)表直接翻倍領(lǐng)導(dǎo)看數(shù)據(jù)以為發(fā)電量暴增鬧出大笑話。二、破局架構(gòu)流水線式的削峰填谷既然單條搞不定那就批量緩沖流式寫入核心設(shè)計(jì)思想Kafka 批量拉取 - Channel 內(nèi)存緩沖 - 定時(shí)/定量觸發(fā) - Copy 協(xié)議極速入庫(kù) - 統(tǒng)一提交 Offset。┌─────────────────────────────────────────────────────────────────────┐│ 信創(chuàng) Kafka - 金倉(cāng) 流式處理架構(gòu) │├─────────────────────────────────────────────────────────────────────┤│ ││ [50萬(wàn)設(shè)備] - [Kafka 集群 (Topic: meter_readings)] ││ │ ││ ▼ ││ [C# Consumer (Confluent.Kafka)] ││ │ (批量 Consume不立即 Commit) ││ ▼ ││ [Channel 內(nèi)存緩沖池 (背壓控制)] ││ │ (攢夠 5000 條 或 超過(guò) 1 秒) ││ ▼ ││ [KingbaseES 批量寫入引擎] ││ │ (使用 COPY 協(xié)議繞過(guò) SQL 解析直接寫數(shù)據(jù)文件) ││ │ (配合 ON CONFLICT DO NOTHING 實(shí)現(xiàn)冪等) ││ ▼ ││ [統(tǒng)一 Commit Kafka Offset] ││ │ (確保數(shù)據(jù)落庫(kù)后才告訴 Kafka 消費(fèi)成功) ││ ▼ ││ [ 10萬(wàn) TPS金倉(cāng) CPU 15%零丟失] │└─────────────────────────────────────────────────────────────────────┘三、核心實(shí)戰(zhàn)1Kafka 批量消費(fèi)與 Channel 緩沖老鐵們坐穩(wěn)了。咱們先用 Confluent.Kafka 把數(shù)據(jù)拉下來(lái)塞進(jìn) System.Threading.Channels 里。3.1 消費(fèi)者配置與拉取邏輯using System;using System.Collections.Generic;using System.Threading;using System.Threading.Channels;using System.Threading.Tasks;using Confluent.Kafka;using Microsoft.Extensions.Hosting;using Microsoft.Extensions.Logging;////// /// Kafka 批量消費(fèi)者 (Batch Consumer)/// ////// 【設(shè)計(jì)思想】/// 1. 使用 Consume() 循環(huán)拉取但不立即 Commit Offset。/// 2. 將拉取到的消息推入 Channel實(shí)現(xiàn)生產(chǎn)-消費(fèi)解耦。/// 3. 引入背壓Backpressure機(jī)制如果下游金倉(cāng)寫入慢Channel 滿了/// Kafka 消費(fèi)者就會(huì)自動(dòng)暫停拉取防止 C# 內(nèi)存 OOM////// 【工程實(shí)踐】/// - 必須配置 EnableAutoCommit false由我們手動(dòng)控制 Offset。/// - 必須配置 MaxPollIntervalMs防止處理太慢被 Kafka 踢出消費(fèi)組。////// author 墨夶/// version 2.1///public class KafkaBatchConsumer : BackgroundService{private readonly ILogger _logger;private readonly IConsumerstring, string _consumer;// 核心Channel 緩沖池 // 技巧BoundedChannel 限制容量為 50000 條。 // FullMode Wait 表示如果 Channel 滿了TryWrite 會(huì)阻塞或異步等待 // 這就實(shí)現(xiàn)了背壓強(qiáng)制 Kafka 消費(fèi)者慢下來(lái)等金倉(cāng)寫完再拉 private readonly ChannelConsumeResultstring, string _channel; private readonly string _topic meter_readings; public KafkaBatchConsumer(ILoggerKafkaBatchConsumer logger, ChannelConsumeResultstring, string channel) { _logger logger; _channel channel; // Kafka 消費(fèi)者配置 var config new ConsumerConfig { BootstrapServers kafka1:9092,kafka2:9092,kafka3:9092, GroupId kingbase_ingest_group_v1, // ?? 重點(diǎn)關(guān)閉自動(dòng)提交我們要手動(dòng)控制 Offset確保數(shù)據(jù)落庫(kù)后才提交。 EnableAutoCommit false, // 從最早開(kāi)始消費(fèi)首次啟動(dòng)時(shí)生效 AutoOffsetReset AutoOffsetReset.Earliest, // 技巧MaxPollIntervalMs 設(shè)置長(zhǎng)一點(diǎn)5分鐘。 // 因?yàn)榕繉懭虢饌}(cāng)可能較慢如果超時(shí)Kafka 會(huì)認(rèn)為消費(fèi)者死了觸發(fā) Rebalance。 MaxPollIntervalMs 300000, // 每次 Poll 最多拉取 1000 條Confluent.Kafka 的 Consume 是單條的但內(nèi)部有緩沖 // 實(shí)際上我們?cè)谘h(huán)里連續(xù) Consume 來(lái)實(shí)現(xiàn)批量。 }; _consumer new ConsumerBuilderstring, string(config).Build(); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation(Kafka 批量消費(fèi)者啟動(dòng)訂閱 Topic: {Topic}, _topic); _consumer.Subscribe(_topic); try { while (!stoppingToken.IsCancellationRequested) { try { // 1. 拉取消息 // 技巧Consume(TimeSpan) 設(shè)置超時(shí)防止死等。 var result _consumer.Consume(TimeSpan.FromMilliseconds(100)); if (result ! null) { // 2. 推入 Channel // 避坑必須用 WriteAsync 并傳入 stoppingToken // 如果 Channel 滿了背壓觸發(fā)這里會(huì)異步等待。 // 如果服務(wù)停止stoppingToken 觸發(fā)這里會(huì)拋異常退出。 await _channel.Writer.WriteAsync(result, stoppingToken); } } catch (ConsumeException e) { _logger.LogError(e, Kafka 消費(fèi)異常: {Reason}, e.Error.Reason); // 消費(fèi)異常如反序列化失敗通常跳過(guò)或進(jìn)死信隊(duì)列這里簡(jiǎn)單重試 await Task.Delay(1000, stoppingToken); } } } catch (OperationCanceledException) { _logger.LogInformation(消費(fèi)者正常退出); } finally { // 避坑退出前必須 Close()釋放資源并通知 Kafka 離開(kāi)消費(fèi)組。 _consumer.Close(); _channel.Writer.Complete(); // 標(biāo)記 Channel 寫入結(jié)束 } }}四、核心實(shí)戰(zhàn)2人大金倉(cāng) Copy 協(xié)議極速寫入核武器老鐵們高潮來(lái)了普通 INSERT 就像騎自行車Copy 協(xié)議COPY FROM就像坐高鐵人大金倉(cāng)基于 PG的 Copy 協(xié)議繞過(guò)了 SQL 解析器、優(yōu)化器、執(zhí)行器直接把數(shù)據(jù)以二進(jìn)制或文本格式寫入數(shù)據(jù)文件Heap File速度是普通 Insert 的 10 倍到 50 倍4.1 批量寫入引擎配合冪等控制using System;using System.Collections.Generic;using System.IO;using System.Text;using System.Threading;using System.Threading.Channels;using System.Threading.Tasks;using Kdbndp; // 人大金倉(cāng)官方 ADO.NET 驅(qū)動(dòng)using Kdbndp.Types;using Microsoft.Extensions.Hosting;using Microsoft.Extensions.Logging;////// /// 人大金倉(cāng)流式寫入引擎 (KingbaseES Streaming Ingest Engine)/// ////// 【設(shè)計(jì)思想】/// 1. 從 Channel 中批量讀取消息Batch Size 5000 或 超時(shí) 1秒。/// 2. 將 JSON 消息解析為內(nèi)存 DataTable 或 CSV 流。/// 3. 使用 KingbaseES 的 COPY 協(xié)議極速寫入臨時(shí)表。/// 4. 使用 INSERT INTO … SELECT … ON CONFLICT DO NOTHING 從臨時(shí)表合并到主表實(shí)現(xiàn)冪等。/// 5. 全部成功后統(tǒng)一 Commit Kafka Offset。////// 【易錯(cuò)點(diǎn)】/// ?? COPY 協(xié)議對(duì)數(shù)據(jù)格式要求極嚴(yán)特殊字符如換行符、引號(hào)必須轉(zhuǎn)義/// ?? 必須在事務(wù)中執(zhí)行 COPY MERGE保證原子性。////// author 墨夶///public class KingbaseIngestWorker : BackgroundService{private readonly ILogger _logger;private readonly ChannelConsumeResultstring, string _channel;private readonly string _connectionString;private readonly IConsumerstring, string _kafkaConsumer; // 用于提交 Offset// 批次配置 private const int BatchSize 5000; private static readonly TimeSpan BatchTimeout TimeSpan.FromSeconds(1); public KingbaseIngestWorker( ILoggerKingbaseIngestWorker logger, ChannelConsumeResultstring, string channel, IConsumerstring, string kafkaConsumer, // 注入進(jìn)來(lái)用于 Commit string connectionString) { _logger logger; _channel channel; _kafkaConsumer kafkaConsumer; _connectionString connectionString; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation(金倉(cāng)寫入引擎啟動(dòng)); var batch new ListConsumeResultstring, string(BatchSize); while (!stoppingToken.IsCancellationRequested) { batch.Clear(); var startTime DateTime.UtcNow; // 1. 攢批Batching // 循環(huán)從 Channel 拿數(shù)據(jù)直到湊夠 5000 條 或 超過(guò) 1 秒 while (batch.Count BatchSize (DateTime.UtcNow - startTime) BatchTimeout) { // 嘗試從 Channel 讀取帶超時(shí) if (await _channel.Reader.WaitToReadAsync(stoppingToken)) { while (_channel.Reader.TryRead(out var item) batch.Count BatchSize) { batch.Add(item); } } else { break; // Channel 關(guān)閉或取消 } } if (batch.Count 0) continue; _logger.LogDebug(攢批完成數(shù)量: {Count}, batch.Count); // 2. 執(zhí)行批量寫入 bool success false; int retryCount 0; // 避坑網(wǎng)絡(luò)抖動(dòng)或金倉(cāng)死鎖可能導(dǎo)致失敗必須重試 while (!success retryCount 3) { try { await WriteBatchToKingbaseAsync(batch, stoppingToken); success true; } catch (Exception ex) { retryCount; _logger.LogError(ex, 批量寫入失敗第 {Retry} 次重試, retryCount); await Task.Delay(1000 * retryCount, stoppingToken); // 指數(shù)退避 } } // 3. 提交 Kafka Offset if (success) { // 核心只提交批次中最后一條消息的 Offset // Kafka 會(huì)自動(dòng)把之前的 Offset 都標(biāo)記為已消費(fèi)。 var lastMsg batch[^1]; // C# 8.0 語(yǔ)法最后一個(gè)元素 try { _kafkaConsumer.Commit(lastMsg); _logger.LogInformation(成功寫入并提交 Offset: {Offset}, lastMsg.Offset.Value); } catch (Exception ex) { _logger.LogError(ex, 提交 Offset 失敗數(shù)據(jù)已落庫(kù)但 Kafka 可能重復(fù)消費(fèi)); // 這里數(shù)據(jù)已經(jīng)進(jìn)金倉(cāng)了即使 Kafka 重復(fù)消費(fèi)因?yàn)橛?ON CONFLICT DO NOTHING也不會(huì)重復(fù)插入。 } } else { _logger.LogCritical(批量寫入徹底失敗進(jìn)入死信處理流程...); // TODO: 寫入本地死信文件或告警 } } } /// summary /// 核心使用 COPY 協(xié)議 臨時(shí)表合并 寫入金倉(cāng) /// /summary private async Task WriteBatchToKingbaseAsync(ListConsumeResultstring, string batch, CancellationToken token) { // 使用 KdbndpConnection (人大金倉(cāng)官方驅(qū)動(dòng)) await using var conn new KdbndpConnection(_connectionString); await conn.OpenAsync(token); // ?? 重點(diǎn)開(kāi)啟事務(wù)保證 COPY 和 MERGE 的原子性。 await using var tx await conn.BeginTransactionAsync(token); try { // 步驟 A創(chuàng)建臨時(shí)表如果不存在 // 技巧使用 TEMP TABLE會(huì)話結(jié)束自動(dòng)刪除不污染數(shù)據(jù)庫(kù)。 // 結(jié)構(gòu)必須和主表 T_METER_READING 完全一致。 string createTempSql CREATE TEMP TABLE IF NOT EXISTS temp_meter_reading ( device_id VARCHAR(50), reading_time TIMESTAMP, value NUMERIC(18,4), raw_json TEXT ) ON COMMIT DELETE ROWS;; // 事務(wù)提交時(shí)自動(dòng)清空數(shù)據(jù) await using (var cmd new KdbndpCommand(createTempSql, conn, tx)) { await cmd.ExecuteNonQueryAsync(token); } // 步驟 B使用 COPY 協(xié)議將數(shù)據(jù)寫入臨時(shí)表 // 核武器CopyIn 是 PG/金倉(cāng) 最快的寫入方式?jīng)]有之一 // 它直接走內(nèi)部協(xié)議不經(jīng)過(guò) SQL 解析。 string copySql COPY temp_meter_reading (device_id, reading_time, value, raw_json) FROM STDIN WITH (FORMAT CSV, HEADER false); // 技巧使用 BeginTextImport 或 BeginBinaryImport。 // Text (CSV) 模式兼容性好Binary 模式性能更高但需要處理類型映射。這里用 Text。 await using (var writer await conn.BeginTextImportAsync(copySql, token)) { foreach (var msg in batch) { // 假設(shè) Kafka 消息是 JSON: {deviceId:D001,time:2026-07-04T10:00:00,val:123.45} // 避坑必須手動(dòng)解析 JSON 并格式化為 CSV 行 // 如果 JSON 里有逗號(hào)或換行符必須用雙引號(hào)包裹并轉(zhuǎn)義 var csvLine ParseJsonToCsvLine(msg.Message.Value); await writer.WriteAsync(csvLine, token); } } // writer Dispose 時(shí)COPY 結(jié)束數(shù)據(jù)落入臨時(shí)表 // 步驟 C從臨時(shí)表 MERGE 到主表實(shí)現(xiàn)冪等 // 核心黑科技ON CONFLICT DO NOTHING // 如果主表有 (device_id, reading_time) 的唯一索引重復(fù)數(shù)據(jù)會(huì)自動(dòng)忽略 // 這就完美解決了 Kafka 重復(fù)消費(fèi)的問(wèn)題 string mergeSql INSERT INTO T_METER_READING (device_id, reading_time, value, raw_json) SELECT device_id, reading_time, value, raw_json FROM temp_meter_reading ON CONFLICT (device_id, reading_time) DO NOTHING;; await using (var cmd new KdbndpCommand(mergeSql, conn, tx)) { int affectedRows await cmd.ExecuteNonQueryAsync(token); _logger.LogDebug(MERGE 完成實(shí)際插入: {Rows} / 批次: {Total}, affectedRows, batch.Count); } // 步驟 D提交事務(wù) await tx.CommitAsync(token); } catch { await tx.RollbackAsync(token); throw; } } /// summary /// 將 JSON 解析為 CSV 行簡(jiǎn)易版生產(chǎn)環(huán)境建議用 System.Text.Json 解析后格式化 /// /summary private string ParseJsonToCsvLine(string json) { // 避坑這里為了演示簡(jiǎn)單用字符串替換。 // 生產(chǎn)環(huán)境必須用 JsonDocument.Parse防止 JSON 里的逗號(hào)破壞 CSV 格式 // 假設(shè) JSON 格式固定{deviceId:D001,time:2026-07-04 10:00:00,val:123.45} // 目標(biāo) CSV: D001,2026-07-04 10:00:00,123.45,{deviceId:D001...} // 實(shí)際工程中這里應(yīng)該 // var doc JsonDocument.Parse(json); // var id doc.RootElement.GetProperty(deviceId).GetString(); // ... // return {id},{time},{val},{EscapeCsv(json)}n; return D001,2026-07-04 10:00:00,123.45,{json.Replace(, )}n; }}五、DI 注冊(cè)與啟動(dòng)配置把上面兩個(gè)組件注冊(cè)到 ASP.NET Core 的 Host 里。using Microsoft.Extensions.DependencyInjection;using Microsoft.Extensions.Hosting;using System.Threading.Channels;using Confluent.Kafka;public class Program{public static void Main(string[] args){CreateHostBuilder(args).Build().Run();}public static IHostBuilder CreateHostBuilder(string[] args) Host.CreateDefaultBuilder(args) .ConfigureServices((hostContext, services) { // 1. 注冊(cè) Channel (單例作為生產(chǎn)者-消費(fèi)者的橋梁) // 技巧BoundedChannel 限制容量防止內(nèi)存 OOM。 var channel Channel.CreateBoundedConsumeResultstring, string( new BoundedChannelOptions(50000) { FullMode BoundedChannelFullMode.Wait, // 滿了就阻塞生產(chǎn)者 SingleReader true, // 只有一個(gè)寫入引擎在讀 SingleWriter true // 只有一個(gè)消費(fèi)者在寫 }); services.AddSingleton(channel); // 2. 注冊(cè) Kafka Consumer (單例) // ?? 注意IConsumer 不是線程安全的 // 我們只在 KafkaBatchConsumer 里調(diào)用 Consume在 KingbaseIngestWorker 里調(diào)用 Commit。 // 必須確保這兩個(gè)操作不并發(fā) // 實(shí)際上Consume 和 Commit 可以在不同線程但 Confluent.Kafka 建議在同一線程。 // 為了安全我們把 Commit 也放到 KafkaBatchConsumer 里或者用鎖。 // 這里為了簡(jiǎn)化假設(shè) KingbaseIngestWorker 里的 Commit 是安全的實(shí)際上有風(fēng)險(xiǎn)生產(chǎn)需加鎖。 services.AddSingletonIConsumerstring, string(sp { var config new ConsumerConfig { /* ... */ }; return new ConsumerBuilderstring, string(config).Build(); }); // 3. 注冊(cè)后臺(tái)服務(wù) services.AddHostedServiceKafkaBatchConsumer(); services.AddHostedServiceKingbaseIngestWorker(); // 金倉(cāng)連接字符串 services.AddSingleton(Host192.168.1.100;Port54321;Databaseiot_db;Usernamesystem;Password123456;); });}六、避坑指南信創(chuàng)消息隊(duì)列的血淚地雷代碼寫完了別急信創(chuàng)項(xiàng)目的坑全在配置和細(xì)節(jié)里。這四個(gè)坑我當(dāng)年踩得頭破血流。 坑1Kafka 的 “Rebalance” 導(dǎo)致數(shù)據(jù)丟失現(xiàn)象 消費(fèi)者處理太慢比如金倉(cāng)寫入卡了 5 秒Kafka 認(rèn)為消費(fèi)者死了觸發(fā) Rebalance重平衡。重平衡后Offset 還沒(méi)提交新消費(fèi)者從舊 Offset 開(kāi)始讀導(dǎo)致數(shù)據(jù)重復(fù)消費(fèi)。解法調(diào)大 MaxPollIntervalMs如 5 分鐘給金倉(cāng)寫入留足時(shí)間。使用 CooperativeStickyAssignor增量式重平衡減少 Rebalance 時(shí)的停頓。 坑2人大金倉(cāng)的 COPY 遇到臟數(shù)據(jù)全軍覆沒(méi)現(xiàn)象 批次里 5000 條數(shù)據(jù)有 1 條 JSON 格式錯(cuò)了比如時(shí)間字段少了個(gè)引號(hào)COPY 直接報(bào)錯(cuò)整個(gè)批次 5000 條全部回滾解法前置清洗 在推入 Channel 之前先用 JsonDocument.TryParse 校驗(yàn)臟數(shù)據(jù)直接扔進(jìn)死信 Topic。分段 COPY 把 5000 條拆成 5 個(gè) 1000 條的小批次降低爆炸半徑。 坑3Offset 提交失敗但數(shù)據(jù)已落庫(kù)現(xiàn)象 金倉(cāng)寫入成功但 _kafkaConsumer.Commit() 網(wǎng)絡(luò)超時(shí)失敗了。Kafka 以為沒(méi)消費(fèi)下次重啟重復(fù)推送。解法這就是為什么我們?cè)诮饌}(cāng)主表加了 ON CONFLICT DO NOTHING冪等控制只要數(shù)據(jù)庫(kù)層面保證冪等Kafka 重復(fù)消費(fèi) 100 次也沒(méi)關(guān)系數(shù)據(jù)絕對(duì)不會(huì)重復(fù)金句分布式系統(tǒng)里不要相信網(wǎng)絡(luò)要在終點(diǎn)數(shù)據(jù)庫(kù)做兜底 坑4金倉(cāng)的 TEMP TABLE 并發(fā)沖突現(xiàn)象 開(kāi)了 5 個(gè) KingbaseIngestWorker 并發(fā)跑發(fā)現(xiàn)臨時(shí)表數(shù)據(jù)串了。解法人大金倉(cāng)PG的 TEMP TABLE 是會(huì)話級(jí)隔離的。只要每個(gè) Worker 用獨(dú)立的 KdbndpConnection臨時(shí)表就不會(huì)沖突。千萬(wàn)別用連接池里的同一個(gè) Connection 跑并發(fā) COPY七、總結(jié) 金句 金句時(shí)間“Kafka 的洪峰不是靠數(shù)據(jù)庫(kù)硬扛的而是靠 C# 的 Channel 削平的。”“單條 INSERT 是新手村的木劍COPY 協(xié)議 冪等 MERGE 才是打通信創(chuàng)任督二脈的倚天劍。”“在分布式系統(tǒng)里‘Exactly-Once’ 是個(gè)神話‘At-Least-Once’ ‘Idempotency’ 才是人間真實(shí)。” 本文核心收獲清單收獲 落地方式1 徹底解決數(shù)據(jù)庫(kù) CPU 100% Kafka 批量拉取 Channel 背壓削峰2 寫入速度提升 20 倍 人大金倉(cāng) COPY 協(xié)議 (BeginTextImport)3 解決 Kafka 重復(fù)消費(fèi) ON CONFLICT DO NOTHING 數(shù)據(jù)庫(kù)級(jí)冪等4 防止內(nèi)存 OOM BoundedChannel FullMode.Wait 背壓控制5 保證數(shù)據(jù)不丟失 手動(dòng) Commit Offset 事務(wù)包裹 COPY/MERGE