资讯动态

Telegraf 集成 Azure Event Hubs 与 IoT Hub:eventhub_consumer 输入插件配置详解与实现原理

发布时间:2026/9/14 16:53:54 来源:尧图企业网站定制
Telegraf 集成 Azure Event Hubs 与 IoT Hubeventhub_consumer 输入插件配置详解与实现原理【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf本文以 Telegraf 仓库中的eventhub_consumer输入插件为对象系统讲解如何将 Azure Event Hubs 与 Azure IoT Hub 中的消息消费为指标数据覆盖 IoT Hub 前置准备、完整的配置参数说明、环境变量鉴权方式、服务型输入插件的行为差异以及基于 tracking metrics 的可靠投递与断点续读原理。读完本文你将能够独立完成 Event Hub 消费链路的配置、理解max_undelivered_messages等关键参数的取舍并掌握该插件从接收消息到输出确认的完整数据流。插件概览从 Azure Event Hubs / IoT Hub 消费消息eventhub_consumer是 Telegraf 的服务型service input输入插件用于从 Azure Event Hubs 与 Azure IoT Hub 实例中持续消费消息并通过可配置的解析器parser将消息体转换为 Telegraf 指标。其核心定位是消息驱动型接入与传统按固定interval轮询采集的插件不同它启动后便监听事件流每收到一条消息即触发解析与指标生成。从仓库元数据看引入版本Telegraf v1.14.0见 plugins/inputs/eventhub_consumer/README.md 头部徽标插件分类iot、messaging平台支持all全平台注册方式在 eventhub_consumer.go 中通过inputs.Add(eventhub_consumer, ...)注册配置段名为[[inputs.eventhub_consumer]]该插件底层基于 Azure 官方的github.com/Azure/azure-event-hubs-go/v3SDK 构建本文后续的源码分析均以 eventhub_consumer.go 为依据。IoT Hub 前置准备三步打通设备消息链路插件的开发主线围绕 Azure IoT Hub 展开官方 README 给出的接入步骤为创建 Azure IoT Hub 实例按 Azure 官方 IoT Hub 文档中的任意指南创建可通过 Azure 门户、CLI 或 ARM 模板完成。创建设备例如一个模拟的 Raspberry Pi 设备在 IoT Hub 的设备管理页面中注册设备并获取设备连接信息。获取连接字符串插件所需的连接字符串位于 IoT Hub 的Shared access policies共享访问策略中其中iothubowner与service两种策略均可用于消费。设备端向 IoT Hub 上报的消息会进入内置的事件中心兼容端点因此 Event Hub 消费者可读取这些消息。说明该插件同时支持通用 Azure Event Hubs非 IoT Hub 场景只需将消息发送到事件中心并从消费组读取即可。两种场景下系统属性System Properties的丰富程度不同IoT Hub 会额外携带IoTHubDeviceConnectionID、IoTHubAuthGenerationID等设备元数据。服务型输入插件的行为差异eventhub_consumer属于服务型输入插件见 docs/includes/service_input.md。普通插件由采集间隔驱动而服务型插件启动一个持续监听的服务等待指标或事件到来。因此它有两个关键差异全局或插件级interval设置可能不生效采集节奏由消息到达速率决定而非定时轮询CLI 选项--test、--test-wait、--once可能不产生输出这些模式面向一次性或按间隔的采集验证无法驱动一个等待事件的服务型插件详细说明可参考 docs/COMMANDS_AND_FLAGS.md。这意味着在部署前验证该插件时不能依赖--test查看样例数据而应直接以守护进程方式运行并通过输出端观察实际流入的消息。配置指南完整参数详解以下为插件官方示例配置源文件见 plugins/inputs/eventhub_consumer/sample.conf与 README 中sample.conf引用一致# Azure Event Hubs service input plugin [[inputs.eventhub_consumer]] ## The default behavior is to create a new Event Hub client from environment variables. ## This requires one of the following sets of environment variables to be set: ## ## 1) Expected Environment Variables: ## - EVENTHUB_CONNECTION_STRING ## ## 2) Expected Environment Variables: ## - EVENTHUB_NAMESPACE ## - EVENTHUB_NAME ## - EVENTHUB_KEY_NAME ## - EVENTHUB_KEY_VALUE ## ## 3) Expected Environment Variables: ## - EVENTHUB_NAMESPACE ## - EVENTHUB_NAME ## - AZURE_TENANT_ID ## - AZURE_CLIENT_ID ## - AZURE_CLIENT_SECRET ## Uncommenting the option below will create an Event Hub client based solely on the connection string. ## This can either be the associated environment variable or hard coded directly. ## If this option is uncommented, environment variables will be ignored. ## Connection string should contain EventHubName (EntityPath) # connection_string ## Set persistence directory to a valid folder to use a file persister instead of an in-memory persister # persistence_dir ## Change the default consumer group # consumer_group ## By default the event hub receives all messages present on the broker, alternative modes can be set below. ## The timestamp should be in RFC 3339 format. ## The 3 options below only apply if no valid persister is read from memory or file (e.g. first run). # from_timestamp # latest true ## Set a custom prefetch count for the receiver(s) # prefetch_count 1000 ## Add an epoch to the receiver(s) # epoch 0 ## Change to set a custom user agent, telegraf is used by default # user_agent telegraf ## To consume from a specific partition, set the partition_ids option. ## An empty array will result in receiving from all partitions. # partition_ids [0,1] ## Max undelivered messages # max_undelivered_messages 1000 ## Set either option below to true to use a system property as timestamp. ## You have the choice between EnqueuedTime and IoTHubEnqueuedTime. ## It is recommended to use this setting when the data itself has no timestamp. # enqueued_time_as_ts true # iot_hub_enqueued_time_as_ts true ## Tags or fields to create from keys present in the application property bag. ## These could for example be set by message enrichments in Azure IoT Hub. # application_property_tags [] # application_property_fields [] ## Tag or field name to use for metadata ## By default all metadata is disabled # sequence_number_field SequenceNumber # enqueued_time_field EnqueuedTime # offset_field Offset # partition_id_tag PartitionID # partition_key_tag PartitionKey # iot_hub_device_connection_id_tag IoTHubDeviceConnectionID # iot_hub_auth_generation_id_tag IoTHubAuthGenerationID # iot_hub_connection_auth_method_tag IoTHubConnectionAuthMethod # iot_hub_connection_module_id_tag IoTHubConnectionModuleID # iot_hub_enqueued_time_field IoTHubEnqueuedTime ## Data format to consume. data_format influx客户端创建与鉴权方式三选一插件默认通过环境变量创建 Event Hub 客户端支持的鉴权组合有三套对应 Azureazure-event-hubs-goSDK 的约定方式环境变量说明1EVENTHUB_CONNECTION_STRING直接使用连接字符串最简单2EVENTHUB_NAMESPACEEVENTHUB_NAMEEVENTHUB_KEY_NAMEEVENTHUB_KEY_VALUE使用共享访问密钥SAS鉴权3EVENTHUB_NAMESPACEEVENTHUB_NAMEAZURE_TENANT_IDAZURE_CLIENT_IDAZURE_CLIENT_SECRET使用 Azure AD 服务主体SPN鉴权如果connection_string被取消注释则完全忽略环境变量仅依据连接字符串创建客户端。注意连接字符串中应包含EntityPath即 EventHubName以便 SDK 定位目标事件中心。这一分支逻辑在 eventhub_consumer.go 中体现ConnectionString非空时调用eventhub.NewHubFromConnectionString(...)否则调用eventhub.NewHubFromEnvironment(...)。断点续读persistence_dirpersistence_dir用于指定 offset 持久化目录。当该值为空时使用内存型 persister重启后从配置的起始位置重新消费当设置为有效目录时则使用文件型 persisterpersist.NewFilePersister见 eventhub_consumer.go将各分区的消费偏移持久化到磁盘实现跨重启的断点续读。对于生产环境建议始终配置该目录以降低重复消费与数据丢失风险。消费起点from_timestamp 与 latest默认行为接收 broker 上当前存在的全部消息from_timestamp从指定时间点开始消费时间格式遵循 RFC 3339offset date-time 格式例如2024-01-01T00:00:00Zlatest true只接收最新消息跳过历史消息。重要前提这两个选项以及默认行为仅在内存/file persister 中读取不到有效 offset即首次运行时生效。一旦已有持久化偏移则一律从上次记录的位置继续这也解释了三个选项只在首次运行时适用的注释。对应实现在configureReceiver()中eventhub_consumer.goFromTimestamp非零则使用ReceiveFromTimestamp否则若Latest为 true 则使用ReceiveWithLatestOffset。接收性能调优prefetch_count默认 1000接收端预取消息条数影响吞吐与内存占用。设置为 0 时未配置不附加该接收选项epoch默认 0为接收器附加 epoch 值。epoch 是 Event Hubs 的独占消费机制——较新的 epoch 接收器会踢掉同一分区上旧 epoch 的接收器用于实现故障转移时快速接管分区user_agent自定义 User-Agent默认为telegraf。源码中若该项为空则使用internal.ProductToken()eventhub_consumer.go即按当前 Telegraf 版本生成的标识。分区控制partition_idspartition_ids用于指定要消费的分区例如[0,1]。空数组表示消费所有分区。源码在Start()中对此做了处理eventhub_consumer.go若未指定分区则通过hub.GetRuntimeInformation(ctx)查询运行时信息获取全部分区列表再为每个分区创建接收器。消息时间戳与元数据处理enqueued_time_as_ts/iot_hub_enqueued_time_as_ts分别使用系统属性EnqueuedTime消息进入 Event Hubs 的时间或IoTHubEnqueuedTime消息进入 IoT Hub 的时间作为指标时间戳。当业务数据本身不含时间戳时强烈建议开启application_property_tags/application_property_fields从消息的 application property bag 中按 key 提取内容分别生成 tag 或 field。典型用途是消费Azure IoT Hub 的 message enrichments消息增强附加的属性和路由信息一系列*_field/*_tag选项为消息的系统属性sequence number、offset、partition id、partition key 以及 IoT Hub 设备连接元数据指定 tag/field 名称。默认全部禁用按需开启。这些选项的具体行为在createMetrics()中有完整映射eventhub_consumer.go例如sequence_number_field写入event.SystemProperties.SequenceNumberenqueued_time_field写入EnqueuedTime的 Unix 毫秒时间戳UnixNano()/int64(time.Millisecond)partition_id_tag写入分区号字符串等。数据格式data_formatdata_format influx指定消息体解析格式默认使用 InfluxDB Line Protocol。Telegraf 支持在data_format中接入任意已注册的解析器可选格式的完整清单见 docs/DATA_FORMATS_INPUT.md包括 JSON、JSON v2、Grok、CSV、Graphite、Prometheus、Value、XPath 等数十种。每种格式有各自的专属配置项选择时需与上游设备/系统实际产出的消息编码对齐。全局配置与插件通用选项与所有插件一样eventhub_consumer支持 Telegraf 的全局与插件通用配置用于修改指标、标签与字段、设置别名以及配置插件顺序详见 docs/CONFIGURATION.md。常用能力包括name_override/name_prefix/name_suffix覆盖或修饰指标名tags为指标附加固定标签interval、metric_batch_size等在插件级覆盖 [agent] 段的全局设置pass/drop、tagpass/tagdrop基于测量名或标签做过滤。由于该插件是服务型输入interval对其采集节奏影响有限但metric_batch_size与flush_interval会直接影响max_undelivered_messages的取值决策见下文。Tracking Metrics 与可靠投递机制该插件支持 tracking metrics见 docs/includes/plugin_tracking_metrics.md。其核心目标是保证数据不丢失Telegraf 会先读取消息并交给输出端在指标成功投递到所有输出之后才向 Event Hub 确认acknowledge该消息。若 Telegraf 中途停止或系统崩溃未完成投递的消息会在恢复后被重新读取。max_undelivered_messages 的取舍max_undelivered_messages默认 1000常量定义见 eventhub_consumer.go限定了已从 broker 读取但尚未被输出端写出的最大消息数本质上是插件内的流量控制窗口设得过高Telegraf 可能持续向输出端推送大批量数据忽视输出端自身的 flush 节奏造成输出端压力与资源占用上升设得过低可能导致 broker 上的消息长期得不到排空消费进度停滞。该值需要与 agent 的metric_batch_size统筹考虑注释原文This value needs to be picked with awareness of the agents metric_batch_size value as well。文档 docs/METRICS.md 也给出了同样的提醒设置过高时Telegraf 可能在每个采集周期都向输出端推送常量批次的指标。源码层面的完整数据流结合 eventhub_consumer.go该机制的实现可以还原为一条清晰的调用链启动Start()L115-L149创建指标通道e.in启动startTracking协程为每个分区调用hub.Receive(ctx, partitionID, e.onMessage, receiveOpts...)注册消息回调接收onMessageL191-L203调用createMetrics将消息体解析为指标写入通道。返回nil表示事件被立即接受并更新 offset返回 error 则标记为重新投递跟踪startTrackingL233-L262将Accumulator包装为带跟踪能力的TrackingAccumulatoracc.WithTracking(e.MaxUndeliveredMessages)用信号量semaphore控制未投递窗口用groups映射保存每个 TrackingID 对应的指标深拷贝副本确认当输出端完成投递后Delivered()通道返回投递信息onDeliveryL206-L231判断track.Delivered()成功则释放信号量槽位失败则利用预先保存的深拷贝副本通过AddTrackingMetricGroup重新加入处理管线注释说明由于onMessage返回时 Event Hub 侧已完成 offset 更新无法依赖 Event Hub 重投故采用本地副本重放。这一设计与 docs/METRICS.md 中先将数据交给输出再向消息源确认的语义完全一致是插件实现不丢消息承诺的根基。指标结构与输出示例README 的 Metrics 与 Example Output 小节当前留空但结合源码可以明确产出的指标形态指标名measurement由data_format对应的解析器决定——以influx为例即消息体 Line Protocol 中的 measurement字段/标签fields/tags来自消息体解析结果叠加application_property_fields/application_property_tags提取的应用属性以及按需开启的各类系统属性字段与 IoT Hub 元数据标签时间戳timestamp默认取消息体自带时间开启enqueued_time_as_ts或iot_hub_enqueued_time_as_ts后分别使用EnqueuedTime/IoTHubEnqueuedTime系统时间。一个典型场景下的指标示意启用sequence_number_field、partition_id_tag、enqueued_time_as_ts时device_temperature,PartitionID3,deviceIdraspberry-pi-1 temperature23.5,SequenceNumber2048i 1735689600000000000实际输出的完整字段集合取决于消息内容与上述元数据选项的开启情况读者可在消费端按需取舍。小结eventhub_consumer是 Telegraf 接入 Azure 消息生态的桥头堡插件具备三个鲜明特点服务型架构事件驱动而非定时轮询、多套 Azure 鉴权方式连接字符串 / SAS / AAD、基于 tracking metrics 的可靠投递。配置时需重点把握三件事一是按部署环境选择鉴权方式并确保EntityPath正确二是通过persistence_dir开启文件持久化以获得断点续读能力三是将max_undelivered_messages与 agent 的metric_batch_size联合调优。本文引用的核心材料包括 插件源码、示例配置、Tracking Metrics 说明 与 输入数据格式手册读者可据此继续深入。【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价