- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本指南以 Apache Pulsar 官方文档《How to develop Pulsar connectors》为骨架,结合本仓库源码(pulsar-io/core、pulsar-functions/api-java、pulsar-io/kafka、pulsar-io/twitter、tests/integration等)展开深入讲解。你将系统掌握 Pulsar Source 与 Sink 连接器的接口契约、Schema 处理(含 KVRecord 与 GenericObject)、单元/集成测试方法、NAR 与 Uber JAR 打包方式,以及如何通过recordMetric为连接器定制监控指标——读完即可上手开发自己的连接器。
连接器是什么
Pulsar Connector 的作用是在 Pulsar 与其他系统之间搬运数据。根据数据流动方向,连接器分为两类:
| 类型 | 说明 | 示例 |
|---|---|---|
| Source | 将数据从外部系统导入 Pulsar | RabbitMQ source connector 把 RabbitMQ 队列中的消息导入 Pulsar topic |
| Sink | 将数据从 Pulsar 导出到外部系统 | Kinesis sink connector 把 Pulsar topic 中的消息导出到 Kinesis stream |
从实现机制上看,Pulsar connector 本质上是特殊的 Pulsar Function,因此开发连接器的方式与开发 Pulsar Function 高度相似:连接器同样运行在 Function 运行时(Function Runtime)之上,由函数工作者(Function Worker)调度、扩缩容与监控。这也意味着函数相关的部署、监控、状态存储等能力对连接器同样适用。
开发 Source 连接器
开发 Source 连接器就是实现 Source 接口。该接口位于pulsar-io/core模块,标注为@InterfaceAudience.Public与@InterfaceStability.Stable,属于稳定的公共 API,仅包含两个方法:
public interface Source<T> extends AutoCloseable { void open(Map<String, Object> config, SourceContext sourceContext) throws Exception; Record<T> read() throws Exception; }实现 open 方法
/** * Open connector with configuration * * @param config initialization config * @param sourceContext * @throws Exception IO type exceptions when opening a connector */ void open(final Map<String, Object> config, SourceContext sourceContext) throws Exception;open方法在 Source 连接器初始化时被调用。在这个方法里:
- 通过传入的
config(Map)读取该连接器所有的自定义配置项,并初始化所需资源(如网络客户端、连接池等); - 将
SourceContext保存下来供后续使用。
例如 Kafka Source 就在open中创建了 Kafka 消费客户端。参考 KafkaAbstractSource.java,其open方法内部做了三件事:
- 用
KafkaSourceConfig.load(config)从config反序列化出强类型的配置对象; - 对
topic、bootstrapServers、groupId等关键配置做Objects.requireNonNull非空校验,并对fetchMinBytes、sessionTimeoutMs、heartbeatIntervalMs等参数做合法性校验; - 将配置拼装成 Kafka
Properties,创建KafkaConsumer并启动拉取线程start()。
这种"config 解析 → 参数校验 → 资源初始化"的三段式写法是内置连接器普遍遵循的模式。
除了获取配置,Pulsar 运行时还通过 SourceContext 向连接器暴露运行环境信息,例如:
getSourceName()/getOutputTopic():当前 Source 的名称与输出 topic;newOutputMessage(topicName, schema):按指定 Schema 构造输出消息;newConsumerBuilder(schema):创建 Pulsar 消费端;- 继承自 BaseContext 的能力:
getLogger()获取日志器、getSecret(name)读取密钥、getStateStore(...)/putState/getState访问状态存储、incrCounter/getCounter使用分布式计数器、recordMetric上报自定义指标、getPulsarClient()获取预配置的 Pulsar 客户端等。
实现 read 方法
/** * Reads the next message from source. * If source does not have any new messages, this call should block. * @return next message from source. The return result should never be null * @throws Exception */ Record<T> read() throws Exception;read方法负责从外部系统读取下一条消息。当没有新消息时,实现应当阻塞等待,而不是返回null——这是该方法最重要的约定。
返回的 Record 需要封装 Pulsar IO 运行时所需的如下信息:
Record 变量
| 变量 | 必填 | 说明 |
|---|---|---|
TopicName | 否 | 该记录来源的 Pulsar topic 名 |
Key | 否 | 消息可选的键。键用于路由,详见 Routing modes |
Value | 是 | 记录的实际数据 |
EventTime | 否 | 记录在源端的事件时间(自 epoch 起的毫秒数) |
PartitionId | 否 | 若记录来自分区源,返回其分区 ID。Pulsar IO 运行时用它与RecordSequence一起构成唯一标识,用于去重并实现 exactly-once 处理保证 |
RecordSequence | 否 | 若记录来自有序源,返回其序号,同样是去重/精确一次语义的组成部分 |
Properties | 否 | 用户自定义属性(Map<String, String>) |
DestinationTopic | 否 | 该消息应被写入的目标 topic,支持按消息粒度路由 |
Message | 否 | 携带用户发送数据的Message对象,见 Message.java |
Record 方法
| 方法 | 说明 |
|---|---|
ack | 确认该记录已被完全处理 |
fail | 标记该记录处理失败 |
Kafka Source 的内部类 KafkaRecord 给出了很好的参考实现:getPartitionId()/getPartitionIndex()返回 Kafka 分区号,getRecordSequence()返回 Kafka offset,getKey()返回消息键,getValue()返回反序列化后的值,ack()通过CompletableFuture完成回调,驱动 Kafka 消费位移提交。
处理 Schema 信息(Source 侧)
Pulsar IO 会自动处理 Schema,并基于 Java 泛型提供强类型 API。如果你明确知道自己要产出的数据类型,可以在 Source 声明中直接指定对应的 Java 类型:
public class MySource implements Source<String> { public Record<String> read() {} }如果你要实现的 Source 需要兼容任意 Schema,可以改用byte[](或ByteBuffer)配合Schema.AUTO_PRODUCE_BYTES():
public class MySource implements Source<byte[]> { public Record<byte[]> read() { Schema wantedSchema = .... Record<byte[]> myRecord = new MyRecordImplementation(); .... } class MyRecordImplementation implements Record<byte[]> { public byte[] getValue() { return ....encoded byte[]...that represents the value } public Schema<byte[]> getSchema() { return Schema.AUTO_PRODUCE_BYTES(wantedSchema); } } }正确处理 KeyValue 类型
要正确处理KeyValue类型,你的 Record 实现需要遵循三条规则:
- 实现 KVRecord 接口,并实现
getKeySchema()、getValueSchema()和getKeyValueEncodingType()三个方法; Record.getValue()必须返回KeyValue对象;Record.getSchema()可以返回null。
当 Pulsar IO 运行时遇到KVRecord时,会自动完成以下转换:
- 正确设置
KeyValueSchema; - 按照
KeyValueEncoding(SEPARATED或INLINE)编码消息键与消息值。
KVRecord<K, V>接口本身位于pulsar-functions/api-java,其getKeyValueEncodingType()返回org.apache.pulsar.common.schema.KeyValueEncodingType枚举。Kafka Source 的 KeyValueKafkaRecord 就是典型实现:它实现了KVRecord<Object, Object>,分别持有 keySchema 与 valueSchema,并返回KeyValueEncodingType.SEPARATED(键值分离编码)。
关于如何实现一个完整的 Source 连接器,可以参考内置的 KafkaSource。
开发 Sink 连接器
开发 Sink 连接器与开发 Source 连接器非常类似:实现 Sink 接口,即实现open与write两个方法:
public interface Sink<T> extends AutoCloseable { void open(Map<String, Object> config, SinkContext sinkContext) throws Exception; void write(Record<T> record) throws Exception; }实现 open 方法
/** * Open connector with configuration * * @param config initialization config * @param sinkContext * @throws Exception IO type exceptions when opening a connector */ void open(final Map<String, Object> config, SinkContext sinkContext) throws Exception;与 Source 的open类似,这里读取配置、初始化外部系统客户端等资源。SinkContext提供 Sink 运行环境,包括:
getSinkName():Sink 名称;getInputTopics():所有输入 topic 列表;getSubscriptionType():订阅类型;seek(topic, partition, messageId):将订阅重置到指定消息 ID;pause(topic, partition)/resume(topic, partition):暂停/恢复消费指定 topic 分区;- 同样继承
BaseContext,可获得日志、密钥、状态存储、计数器、指标上报等能力。
实现 write 方法
/** * Write a message to Sink * @param record record to write to sink * @throws Exception */ void write(Record<T> record) throws Exception;在write的实现中,你可以自行决定如何把Value和Key写入目标系统,并利用PartitionId、RecordSequence等信息实现不同的处理保证(例如精确一次语义)。
此外,你必须负责 ack 与 fail 的调用:消息成功写出后调用record.ack(),发送失败则调用record.fail()。这是连接器与运行时之间关于消息处理状态的关键握手——ack/fail 结果会驱动 Pulsar 侧的消费位点推进与重试策略。
处理 Schema 信息(Sink 侧)
与 Source 相同,Pulsar IO 自动处理 Schema,并基于 Java 泛型提供强类型 API。若已知消费的数据类型,可在 Sink 声明中直接指定:
public class MySink implements Sink<String> { public void write(Record<String> record) {} }若需要实现可适配任意 Schema 的 Sink,可以使用特殊的GenericObject接口:
public class MySink implements Sink<GenericObject> { public void write(Record<GenericObject> record) { Schema schema = record.getSchema(); GenericObject genericObject = record.getValue(); if (genericObject != null) { SchemaType type = genericObject.getSchemaType(); Object nativeObject = genericObject.getNativeObject(); ... } .... } }对于 AVRO、JSON 和 Protobuf 类型的记录(schemaType为AVRO、JSON、PROTOBUF_NATIVE),可以把genericObject强转为GenericRecord,使用getFields()和getField()API 访问字段;也可以通过genericObject.getNativeObject()拿到原生 AVRO 记录。
对于KeyValue类型,可以同时访问键 Schema 与值 Schema:
public class MySink implements Sink<GenericObject> { public void write(Record<GenericObject> record) { Schema schema = record.getSchema(); GenericObject genericObject = record.getValue(); SchemaType type = genericObject.getSchemaType(); Object nativeObject = genericObject.getNativeObject(); if (type == SchemaType.KEY_VALUE) { KeyValue keyValue = (KeyValue) nativeObject; Object key = keyValue.getKey(); Object value = keyValue.getValue(); KeyValueSchema keyValueSchema = (KeyValueSchema) schema; Schema keySchema = keyValueSchema.getKeySchema(); Schema valueSchema = keyValueSchema.getValueSchema(); } .... } }关于GenericObject相关能力的验证,可以参考集成测试中的 PulsarGenericObjectSinkTest.java,它覆盖了 Sink 以GenericObject消费不同 Schema 数据的端到端路径。
测试连接器
测试连接器颇具挑战性,因为 Pulsar IO 连接器同时与两个系统交互——Pulsar 本身以及它连接的外部系统,两者都不容易 mock。官方推荐的做法是:在 mock 外部服务的前提下,按下面两级结构编写测试。
单元测试(Unit test)
为连接器创建单元测试,聚焦连接器内部逻辑:配置解析、Record 封装、数据转换、ack/fail 行为等。由于不依赖真实的外部系统与 Pulsar 集群,单元测试应当快速、稳定、覆盖面广。
集成测试(Integration test)
在单元测试足够充分之后,增加独立的集成测试来验证端到端功能。Pulsar 项目全部集成测试均使用 testcontainers——通过容器化方式拉起真实的外部系统(如 RabbitMQ、Kafka),配合测试中的 Pulsar 集群,完成"外部系统 → Source → Pulsar → Sink → 外部系统"的完整闭环验证。
仓库中的参考实现位于 tests/integration/src/test/java/org/apache/pulsar/tests/integration/io,其中:
- PulsarIOTestBase.java 与 PulsarIOTestRunner.java 提供了 IO 集成测试的基类与运行入口;
- RabbitMQSourceTester.java 与 RabbitMQSinkTester.java 示范了如何用 testcontainers 拉起 RabbitMQ 容器并验证 Source/Sink 的数据通路;
sources/与sinks/子目录下按连接器类型组织更多 tester 类。
关于如何为 Pulsar 连接器编写集成测试,可深入参考 tests/integration/src/test/java/org/apache/pulsar/tests/integration/io 中的源码。
打包连接器
开发和测试完成后,需要把连接器打包,才能提交到 Pulsar Functions 集群上运行。与 Function 运行时协作有NAR与Uber JAR两种方式。
许可证与版权提醒:如果你打算把连接器打包分发给他人的话,你有义务妥善处理许可证与版权——为你代码使用的所有库以及你的发行物添加许可证与版权声明。若采用 NAR 方式,NAR 插件会在生成的 NAR 包中自动生成
DEPENDENCIES文件,包含连接器所有库的正确许可证与版权信息。
方式一:NAR(推荐)
NAR(NiFi Archive)是 Apache NiFi 使用的一种自定义打包机制,用于提供一定程度的 Java ClassLoader 隔离。Pulsar 用同样的机制打包全部内置连接器(见 pulsar-io 下各连接器模块)。
打包连接器最简单的方式是使用 nifi-nar-maven-plugin 创建 NAR 包。在连接器模块的 Maven 工程中加入该插件:
<plugins> <plugin> <groupId>org.apache.nifi</groupId> <artifactId>nifi-nar-maven-plugin</artifactId> <version>1.2.0</version> </plugin> </plugins>同时必须在resources/META-INF/services/pulsar-io.yaml中创建如下内容的文件:
name: connector name description: connector description sourceClass: fully qualified class name (only if source connector) sinkClass: fully qualified class name (only if sink connector)这个 YAML 是 NAR 包被 Pulsar 识别为连接器的关键元数据。仓库中每个内置连接器都带有该文件,例如 RabbitMQ 连接器的 pulsar-io.yaml:
name: rabbitmq description: RabbitMQ source and sink connector sourceClass: org.apache.pulsar.io.rabbitmq.RabbitMQSource sinkClass: org.apache.pulsar.io.rabbitmq.RabbitMQSink sourceConfigClass: org.apache.pulsar.io.rabbitmq.RabbitMQSourceConfig sinkConfigClass: org.apache.pulsar.io.rabbitmq.RabbitMQSinkConfig可见实际内置连接器还会补充sourceConfigClass与sinkConfigClass字段,用于声明连接器的配置类,提交连接器时即可通过 YAML 文件直接填充配置。
对于 Gradle 用户,Gradle Plugin Portal 上提供了对应的 Gradle Nar 插件(io.github.lhotari.gradle-nar-plugin)。
关于如何使用 NAR 打包 Pulsar 连接器,可以参考内置连接器的完整工程配置 TwitterFirehose 的 pom.xml——该模块在
<build><plugins>中仅声明nifi-nar-maven-plugin,并依赖pulsar-io-core与pulsar-io-common等模块,即可产出可被 Pulsar 加载的 NAR 包。
方式二:Uber JAR
另一种做法是创建包含连接器全部 JAR 文件及其他资源文件的uber JAR,无需任何目录内部结构。使用 maven-shade-plugin 构建:
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.1.1</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <filters> <filter> <artifact>*:*</artifact> </filter> </filters> </configuration> </execution> </executions> </plugin>两种方式各有取舍:NAR 提供 ClassLoader 隔离、由插件自动生成依赖许可文件,是 Pulsar 内置连接器与社区分发的主流选择;Uber JAR 结构简单、无需额外插件产物,适合快速分发,但需要注意依赖冲突与许可证声明的完整性问题。
监控连接器
Pulsar 连接器让数据进出 Pulsar 变得容易,但确保运行中的连接器时刻健康同样重要。可以通过以下方式监控已部署的连接器:
查看 Pulsar 提供的指标:Pulsar 连接器会暴露指标,可用于监控Java连接器的健康状态。具体采集方式参考 监控指南。
设置并查看自定义指标:除 Pulsar 自带指标外,Pulsar 允许为Java连接器定制指标。函数工作者(Function Worker)会自动把用户自定义指标收集到 Prometheus,并可在 Grafana 中查看。
自定义 Java 连接器指标的示例如下:
public class TestMetricSink implements Sink<String> { @Override public void open(Map<String, Object> config, SinkContext sinkContext) throws Exception { sinkContext.recordMetric("foo", 1); } @Override public void write(Record<String> record) throws Exception { } @Override public void close() throws Exception { } }recordMetric(String metricName, double value)定义在 BaseContext 中,Source 与 Sink 的 Context 均继承该方法,因此在open、write/read的任何位置都能上报自定义指标。指标会被 Function Worker 自动汇总并暴露给 Prometheus 抓取,配合 Grafana 即可构建连接器专属的监控面板。
总结
开发 Pulsar 连接器的完整流程可以概括为:
- 实现接口:Source 实现
open+read,Sink 实现open+write,核心是正确封装 Record(Source 侧)并在write中妥善调用ack/fail(Sink 侧); - 处理 Schema:按需选择强类型泛型、
byte[]+Schema.AUTO_PRODUCE_BYTES、KVRecord(键值场景)或GenericObject(Sink 通用场景); - 分层测试:先单元测试,再用 testcontainers 编写端到端集成测试(参考
tests/integration); - 打包分发:优先 NAR(含
pulsar-io.yaml元数据),或选择 Uber JAR,并处理好许可证声明; - 上线监控:利用 Pulsar 指标与
recordMetric自定义指标,经 Prometheus/Grafana 保障连接器健康。
从仓库源码可以看到,Pulsar 内置连接器(Kafka、RabbitMQ、Kinesis 等)全部遵循同一套接口与打包规范,因此掌握本指南的 API 契约与工程流程后,你既可以复刻内置连接器的成熟模式,也可以自由开发适配任意外部系统的自定义连接器。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar Connector 开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控
Apache Pulsar Connector 开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控 本指南以 Apache Pulsar
消息队列后端流处理Apache Pulsar 内置连接器(Built-in Connector)完全指南:Source 与 Sink 生态全览及实战部署
Apache Pulsar 内置连接器(Built in Connector)完全指南:Source 与 Sink 生态全览及实战部署 Pulsar 发行版中打
消息队列后端流处理Apache Pulsar 内置连接器完全指南:Source 与 Sink 全清单、配置与实战
Apache Pulsar 内置连接器完全指南:Source 与 Sink 全清单、配置与实战 Apache Pulsar 发行版内置了一组经过打包与联调验证的
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考