ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Pulsar InfluxDB Sink Connector 使用指南:从配置解析到批量写入

Apache Pulsar InfluxDB Sink Connector 使用指南:从配置解析到批量写入 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 内置了丰富的 IO 连接器Source/Sink其中InfluxDB Sink Connector用于将 Pulsar 主题中的消息持续拉取并写入 InfluxDB 时序数据库实现流式数据到时序指标的落库。本文以版本化文档 io-influxdb.md 为骨架结合pulsar-io/influxdb模块源码与测试完整讲解该连接器的全部配置项、数据点构建规则、批量刷写机制以及部署运行命令帮助你直接落地一套可用的 InfluxDB 数据管道。InfluxDB Sink 是什么按照 io-overview.md 的定义Pulsar IO 连接器分为两类Source 负责把外部系统数据拉入PulsarSink 负责把 Pulsar 主题数据写出到外部系统。InfluxDB Sink 正属于后者The InfluxDB Sink Connector is used to pull messages from Pulsar topics and persist the messages to an InfluxDB database.它在 io-connectors.md 的内置 Sink 列表中被收录部署时只需指定--sink-type influxdb即可。从仓库源码看该连接器位于 pulsar-io/influxdb 模块Maven 坐标pulsar-io-influxdb其pom.xml同时依赖两代 InfluxDB Java 客户端org.influxdb:influxdb-java:2.22InfluxDB 1.x 客户端v1包com.influxdb:influxdb-client-java:4.0.0InfluxDB 2.x 客户端v2包。因此当前仓库同时提供v1用户名/密码 database与v2token organization bucket两套实现下文先以版本 2.3.1 文档所描述的v1 配置为主再补充 v2 差异。Sink 配置项详解v1原文档给出的配置表如下其中influxdbUrl与database为必填项NameDefaultRequiredDescriptioninfluxdbUrlnulltrueThe url of the InfluxDB instance to connect to.usernamenullfalseThe username used to authenticate to InfluxDB.passwordnullfalseThe password used to authenticate to InfluxDB.databasenulltrueThe InfluxDB database to write to.consistencyLevelONEfalseThe consistency level for writing data to InfluxDB. Possible values [ALL, ANY, ONE, QUORUM].logLevelNONEfalseThe log level for InfluxDB request and response. Possible values [NONE, BASIC, HEADERS, FULL].retentionPolicyautogenfalseThe retention policy for the InfluxDB database.gzipEnablefalsefalseFlag to determine if gzip should be enabled.batchTimeMs1000falseThe InfluxDB operation time in milliseconds.batchSize200falseThe batch size of write to InfluxDB database.上述字段在源码 v1/InfluxDBSinkConfig.java 中一一对应并带FieldDoc注解标注必填性、默认值与说明。其validate()方法L118-L123会校验influxdbUrl、database非空缺失时抛出 property not set 异常batchSize 0、batchTimeMs 0。各参数在底层如何生效结合 InfluxDBBuilderImpl.java可以看到参数的底层作用认证方式当username非空时调用InfluxDBFactory.connect(url, username, password)启用认证否则使用无认证的InfluxDBFactory.connect(url)gzipEnable为true时调用influxDB.enableGzip()开启请求压缩适合高吞吐写入场景以降低网络带宽logLevel通过InfluxDB.LogLevel.valueOf(...)解析为NONE / BASIC / HEADERS / FULL非法值会抛出带合法取值列表的IllegalArgumentException。consistencyLevel的解析位于 InfluxDBAbstractSink.java使用InfluxDB.ConsistencyLevel.valueOf(...)将字符串转为枚举再在批量写入时通过BatchPoints.Builder.consistency(...)应用到每一次写入请求。此外Sink 在open()阶段L64-L68会调用influxDB.describeDatabases()检查目标库是否存在不存在则自动createDatabase(influxDatabase)——也就是说database指向的库无需预先手工创建。完整配置示例YAML连接器配置支持从 YAML 文件或 Map 加载两种方式分别对应 InfluxDBSinkConfig.java 中的load(String yamlFile)与load(MapString, Object)。一份完整的 v1 Sink 配置如下influxdbUrl: http://localhost:8086 # 必填InfluxDB 服务地址 username: admin # 可选认证用户名非空时启用认证 password: password # 可选认证密码 database: test_db # 必填目标数据库不存在则自动创建 consistencyLevel: ONE # 可选默认 ONE取值 [ALL, ANY, ONE, QUORUM] logLevel: NONE # 可选默认 NONE取值 [NONE, BASIC, HEADERS, FULL] retentionPolicy: autogen # 可选默认 autogen gzipEnable: false # 可选是否启用 gzip 压缩 batchTimeMs: 1000 # 可选批量刷写周期毫秒 batchSize: 200 # 可选单批最大消息数该示例与单元测试 v1/InfluxDBSinkConfigTest.java 中loadFromYamlFileTest/loadFromMapTest的断言一致可据此验证配置解析正确性。部署运行创建 InfluxDB Sink按照 io-managing.md 的说明作为内置 Sink提交时无需指定--classname与--archive只要给出--sink-type influxdb./bin/pulsar-admin sinks create \ --tenant tenant \ --namespace namespace \ --name influxdb-sink \ --inputs input-topics \ --sink-type influxdb \ --sink-config-file path-to-influxdb-sink-config.yaml如果希望先在本地以独立进程方式调试运行可使用localrun./bin/pulsar-admin sinks localrun \ --tenant tenant \ --namespace namespace \ --name influxdb-sink \ --inputs input-topics \ --sink-type influxdb \ --sink-config-file path-to-influxdb-sink-config.yaml部署完成后可通过bin/pulsar-admin functions get --tenant ... --namespace ... --name influxdb-sink获取连接器元数据与运行状态连接器本质上是运行在 Pulsar Functions 框架上的实例。注意内置连接器的sink-type取值由pulsar-io.yaml中声明的name决定InfluxDB 即influxdb。消息数据格式如何把 Pulsar 消息变成 InfluxDB Pointv1 Sink 要求输入消息为GenericRecord带 Schema 的记录由 InfluxDBGenericRecordSink.java 负责把每条记录转换成 InfluxDB 的Point。转换规则如下measurement必填记录必须包含名为measurement的字段作为时序数据的 measurement 名缺失时抛出SchemaSerializationException(measurement is a required field.)tags可选若记录包含名为tags的字段且其值为Map类型则键值对全部作为 Point 的 tag非 Map 类型或缺失时忽略tag 为空时间戳默认取System.currentTimeMillis()毫秒精度普通字段除measurement、tags两个保留字段外其余字段如model、value全部作为 Point 的 field 写入。以测试 v1/InfluxDBGenericRecordSinkTest.java 中的Cpu记录为例一条形如{ measurement: cpu, model: lenovo, value: 10, tags: { host: server-1 } }的消息会被转换成 measurement 为cpu、tag 为hostserver-1、field 为model与value的 Point。批量写入机制BatchSink 的刷写与确认语义InfluxDB Sink 并非逐条写入而是继承抽象基类 BatchSink.java 实现批量写入其内部逻辑是init(batchTimeMs, batchSize)创建单线程调度器每隔batchTimeMs毫秒触发一次flush()write(record)将消息加入内存缓冲列表当累积数量达到batchSize时立即提交一次flush()flush()将缓冲列表整体取出逐条调用抽象方法buildPoint(record)转换为 Point转换失败的单条消息调用record.fail()并从列表移除所有 Point 组装完成后调用writePoints(points)v1 实现见 InfluxDBAbstractSink.java通过BatchPoints指定 database、retentionPolicy、consistency 后一次性写入写入成功则对整批消息record.ack()写入失败则整批record.fail()并记录错误日志。因此batchTimeMs与batchSize是时间触发 数量触发的双阈值实际刷写发生在两者先满足其一之时。这决定了 Sink 的端到端延迟单条消息最长会在缓冲区内停留约batchTimeMs毫秒。若你的场景需要更低延迟可适当调小batchTimeMs若追求吞吐可调大batchSize。补充v2InfluxDB 2.x实现差异当前仓库在v1之外还提供了面向 InfluxDB 2.x 的实现v2/InfluxDBSinkConfig.java 与 v2/InfluxDBSink.java主要差异包括认证与目标模型用token必填、标记为敏感、organization必填、bucket必填替代 v1 的username/password/database并通过 InfluxDBClientBuilderImpl.java 构造InfluxDBClientOptions客户端新增precision时间戳精度可配置为ns / us / ms / s默认ns写入时通过WritePrecision.fromValue(...)解析时间戳来源buildPoint优先读取记录中的timestamp字段支持 Number 或 String 类型缺失时回退为当前系统时间字段组织方式记录结构要求measurement、tags、fields三个顶层字段其中tags与fields支持GenericRecordJSON Schema或MapAvro Schema两种形态field 值仅接受 Number、Boolean、String 与 AvroUtf8类型。从源码结构看v2 是面向新版本 InfluxDB 的演进实现2.3.1 版本文档中的配置表对应 v1实际使用前请确认目标 InfluxDB 的版本再选择对应的配置形态。结语InfluxDB Sink 连接器让Pulsar 消息 → 时序数据库的链路只需一份 YAML 配置即可打通influxdbUrl与database决定写入目标consistencyLevel/retentionPolicy/gzipEnable控制写入行为batchTimeMs与batchSize决定批量刷写的延迟与吞吐。配合measurement、tags与普通字段的记录映射规则即可把带 Schema 的 Pulsar 消息稳定落库。若需进一步了解连接器的整体概念与部署管理可继续阅读 io-overview.md 与 io-managing.md底层实现与测试用例可在 pulsar-io/influxdb 模块中深入查阅。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐SeaTunnel Neo4j Sink Connector 使用指南从逐条写入到 UNWIND 批量写入SeaTunnel Neo4j Sink Connector 使用指南从逐条写入到 UNWIND 批量写入 SeaTunnel 的 Neo4j Sink 插件数据工程大数据批处理流处理Apache Pulsar HBase Sink Connector将 Topic 消息批量写入 HBase 表的配置与原理Apache Pulsar HBase Sink Connector将 Topic 消息批量写入 HBase 表的配置与原理 本篇围绕 Pulsar IO 的消息队列后端流处理Apache Pulsar HDFS Sink Connector 实战指南从 Pulsar Topic 写入 HDFS 的配置、原理与源码解析Apache Pulsar HDFS Sink Connector 实战指南从 Pulsar Topic 写入 HDFS 的配置、原理与源码解析 Apache消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表