☰
Apache Pulsar InfluxDB Sink Connector 使用指南:从配置解析到批量写入
2026/9/29 20:29:17 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache 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 负责把外部系统数据拉入Pulsar,Sink 负责把 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.22:InfluxDB 1.x 客户端(v1包);
  • com.influxdb:influxdb-client-java:4.0.0:InfluxDB 2.x 客户端(v2包)。

因此当前仓库同时提供v1(用户名/密码 + database)与v2(token + organization + bucket)两套实现,下文先以版本 2.3.1 文档所描述的v1 配置为主,再补充 v2 差异。

Sink 配置项详解(v1)

原文档给出的配置表如下,其中influxdbUrl与database为必填项:

NameDefaultRequiredDescription
influxdbUrlnulltrueThe 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(Map<String, 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 Point

v1 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 为host=server-1、field 为model与value的 Point。

批量写入机制:BatchSink 的刷写与确认语义

InfluxDB Sink 并非逐条写入,而是继承抽象基类 BatchSink.java 实现批量写入,其内部逻辑是:

  1. init(batchTimeMs, batchSize)创建单线程调度器,每隔batchTimeMs毫秒触发一次flush();
  2. write(record)将消息加入内存缓冲列表,当累积数量达到batchSize时立即提交一次flush();
  3. flush()将缓冲列表整体取出,逐条调用抽象方法buildPoint(record)转换为 Point;转换失败的单条消息调用record.fail()并从列表移除;
  4. 所有 Point 组装完成后调用writePoints(points)(v1 实现见 InfluxDBAbstractSink.java:通过BatchPoints指定 database、retentionPolicy、consistency 后一次性写入);
  5. 写入成功则对整批消息record.ack();写入失败则整批record.fail()并记录错误日志。

因此,batchTimeMs与batchSize是"时间触发 + 数量触发"的双阈值,实际刷写发生在两者先满足其一之时。这决定了 Sink 的端到端延迟:单条消息最长会在缓冲区内停留约batchTimeMs毫秒。若你的场景需要更低延迟,可适当调小batchTimeMs;若追求吞吐,可调大batchSize。

补充:v2(InfluxDB 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支持GenericRecord(JSON Schema)或Map(Avro 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 模块中深入查阅。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询