- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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为必填项:
| Name | Default | Required | Description |
|---|---|---|---|
influxdbUrl | null | true | The url of the InfluxDB instance to connect to. |
username | null | false | The username used to authenticate to InfluxDB. |
password | null | false | The password used to authenticate to InfluxDB. |
database | null | true | The InfluxDB database to write to. |
consistencyLevel | ONE | false | The consistency level for writing data to InfluxDB. Possible values [ALL, ANY, ONE, QUORUM]. |
logLevel | NONE | false | The log level for InfluxDB request and response. Possible values [NONE, BASIC, HEADERS, FULL]. |
retentionPolicy | autogen | false | The retention policy for the InfluxDB database. |
gzipEnable | false | false | Flag to determine if gzip should be enabled. |
batchTimeMs | 1000 | false | The InfluxDB operation time in milliseconds. |
batchSize | 200 | false | The 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 实现批量写入,其内部逻辑是:
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。
补充: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
相关推荐
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),仅供参考