You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

关于使用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
}

目前我们有两个核心疑问:

  1. 该方法每秒会被调用300-500次,已知客户端内置批处理功能,是否需要采用单例模式的客户端?
  2. 当前每个事件创建一个流再调用摄入的方式是否合理,能否配置客户端后直接将事件入队处理?

问题解答

1. 是否需要单例模式的客户端?

必须采用单例模式。Kusto的QueuedIngestClient内部维护了批处理队列、连接池和重试机制,每次调用SaveEvent都创建新客户端会引发以下问题:

  • 完全丧失内置批处理能力:每个客户端的批处理队列独立,无法将多个事件合并为一个批次发送,导致大量小批次请求,既降低摄入效率,也容易触发ADX的请求频率限制。
  • 资源严重浪费:频繁创建销毁客户端会重复初始化连接、认证会话,在每秒300-500次调用的场景下,CPU和网络开销会急剧上升。
  • 配置冗余:每次都要重新设置重试策略、摄入属性等,代码冗余且易出错。

正确做法是将ingestClient作为单例(比如通过依赖注入容器注册为单例实例),在应用启动时初始化一次,后续所有SaveEvent调用复用同一个客户端实例。

2. 单个事件创建流的方式是否合理?能否直接入队事件?

当前逐个创建流的方式功能可行,但存在优化空间;Kusto客户端没有直接入队.NET对象的API,但可以简化序列化流程:

  • 现有方式的问题:每个事件都创建MemoryStream、StreamWriter、JsonTextWriter,存在一定的对象创建开销。
  • 优化建议:
    1. 复用JsonSerializer实例:将JsonSerializer设为单例,避免每次创建新实例的开销。
    2. 简化序列化流程:使用更高效的序列化方式(如System.Text.Json替代Newtonsoft.Json),减少中间对象分配。
    3. 依赖客户端批处理:单例客户端会自动根据数据大小、时间间隔触发批处理(可通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 02:05:25