关于使用Azure Data Explorer SDK摄入流数据的技术疑问
Azure Data Explorer 事件摄入优化问题
我们需要开发一款可逐个接收事件并将其摄入Azure Data Explorer(ADX)的软件,但对Kusto Client的用法存在困惑。以下是我们实现的SaveEvent方法代码:
public void SaveEvent(Object @event) { var _kcsb = new KustoConnectionStringBuilder(Uri).WithAadApplicationKeyAuthentication( applicationClientId: "", applicationKey: "", authority: TenantId); using var ingestClient = KustoIngestFactory.CreateQueuedIngestClient(_kcsb); //// Create your custom implementation of IRetryPolicy, which will affect how the ingest client handles retrying on transient failures IRetryPolicy retryPolicy = new NoRetry(); //// This line sets the retry policy on the ingest client that will be enforced on every ingest call from here on ((IKustoQueuedIngestClient)ingestClient).QueueOptions.QueueRequestOptions.RetryPolicy = retryPolicy; var ingestProperties = new KustoIngestionProperties(DatabaseName, TableName) { Format = DataSourceFormat.json, IngestionMapping = new IngestionMapping { IngestionMappingKind = Kusto.Data.Ingestion.IngestionMappingKind.Json, IngestionMappingReference = MappingName } }; // Build the stream var stream = new MemoryStream(); using var streamWriter = new StreamWriter(stream: stream, encoding: Encoding.UTF8, bufferSize: 4096, leaveOpen: true); using var jsonWriter = new JsonTextWriter(streamWriter); packet.Id = DateTime.UtcNow.Ticks; // 原代码中packet变量未定义,推测为事件对象的属性 var serializer = new JsonSerializer(); serializer.Serialize(jsonWriter, @event); streamWriter.Flush(); stream.Seek(0, SeekOrigin.Begin); // Tell the client to ingest this await ingestClient.IngestFromStreamAsync(stream, ingestProperties); // 修正原代码中参数笔误:data改为stream }
目前我们有两个核心疑问:
- 该方法每秒会被调用300-500次,已知客户端内置批处理功能,是否需要采用单例模式的客户端?
- 当前每个事件创建一个流再调用摄入的方式是否合理,能否配置客户端后直接将事件入队处理?
问题解答
1. 是否需要单例模式的客户端?
必须采用单例模式。Kusto的QueuedIngestClient内部维护了批处理队列、连接池和重试机制,每次调用SaveEvent都创建新客户端会引发以下问题:
- 完全丧失内置批处理能力:每个客户端的批处理队列独立,无法将多个事件合并为一个批次发送,导致大量小批次请求,既降低摄入效率,也容易触发ADX的请求频率限制。
- 资源严重浪费:频繁创建销毁客户端会重复初始化连接、认证会话,在每秒300-500次调用的场景下,CPU和网络开销会急剧上升。
- 配置冗余:每次都要重新设置重试策略、摄入属性等,代码冗余且易出错。
正确做法是将ingestClient作为单例(比如通过依赖注入容器注册为单例实例),在应用启动时初始化一次,后续所有SaveEvent调用复用同一个客户端实例。
2. 单个事件创建流的方式是否合理?能否直接入队事件?
当前逐个创建流的方式功能可行,但存在优化空间;Kusto客户端没有直接入队.NET对象的API,但可以简化序列化流程:
- 现有方式的问题:每个事件都创建
MemoryStream、StreamWriter、JsonTextWriter,存在一定的对象创建开销。 - 优化建议:
- 复用
JsonSerializer实例:将JsonSerializer设为单例,避免每次创建新实例的开销。 - 简化序列化流程:使用更高效的序列化方式(如
System.Text.Json替代Newtonsoft.Json),减少中间对象分配。 - 依赖客户端批处理:单例客户端会自动根据数据大小、时间间隔触发批处理(可通过
QueueOptions调整阈值),无需手动合并事件。
- 复用
示例优化后的代码片段(单例客户端+简化序列化):
// 单例客户端,应用启动时初始化 private static readonly IKustoQueuedIngestClient _singletonIngestClient; private static readonly JsonSerializer _jsonSerializer = new JsonSerializer(); static YourClassName() { var _kcsb = new KustoConnectionStringBuilder(Uri).WithAadApplicationKeyAuthentication( applicationClientId: "", applicationKey: "", authority: TenantId); _singletonIngestClient = (IKustoQueuedIngestClient)KustoIngestFactory.CreateQueuedIngestClient(_kcsb); _singletonIngestClient.QueueOptions.QueueRequestOptions.RetryPolicy = new NoRetry(); } public async Task SaveEvent(Object @event) { var ingestProperties = new KustoIngestionProperties(DatabaseName, TableName) { Format = DataSourceFormat.json, IngestionMapping = new IngestionMapping { IngestionMappingKind = IngestionMappingKind.Json, IngestionMappingReference = MappingName } }; // 简化序列化流程 using var stream = new MemoryStream(); using var streamWriter = new StreamWriter(stream, Encoding.UTF8, 4096, leaveOpen: true); using var jsonWriter = new JsonTextWriter(streamWriter); // 为事件设置Id(根据实际类型调整) ((dynamic)@event).Id = DateTime.UtcNow.Ticks; _jsonSerializer.Serialize(jsonWriter, @event); streamWriter.Flush(); stream.Seek(0, SeekOrigin.Begin); await _singletonIngestClient.IngestFromStreamAsync(stream, ingestProperties); }
若想进一步降低开销,可提前将事件序列化为JSON字符串或字节数组后再传入IngestFromStreamAsync,无需每次创建流对象。
内容的提问来源于stack exchange,提问作者Zapnologica
相关产品推荐
相关产品推荐

