Apache Pulsar 连接器(Connector)开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控
2026/9/24 5:25:19 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

本指南以 Apache Pulsar 官方文档《How to develop Pulsar connectors》为骨架,结合本仓库源码(pulsar-io/corepulsar-functions/api-javapulsar-io/kafkapulsar-io/twittertests/integration等)展开深入讲解。你将系统掌握 Pulsar Source 与 Sink 连接器的接口契约、Schema 处理(含 KVRecord 与 GenericObject)、单元/集成测试方法、NAR 与 Uber JAR 打包方式,以及如何通过recordMetric为连接器定制监控指标——读完即可上手开发自己的连接器。

连接器是什么

Pulsar Connector 的作用是在 Pulsar 与其他系统之间搬运数据。根据数据流动方向,连接器分为两类:

类型说明示例
Source将数据从外部系统导入 PulsarRabbitMQ 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方法内部做了三件事:

  1. KafkaSourceConfig.load(config)config反序列化出强类型的配置对象;
  2. topicbootstrapServersgroupId等关键配置做Objects.requireNonNull非空校验,并对fetchMinBytessessionTimeoutMsheartbeatIntervalMs等参数做合法性校验;
  3. 将配置拼装成 KafkaProperties,创建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
  • 按照KeyValueEncodingSEPARATEDINLINE)编码消息键与消息值。

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 接口,即实现openwrite两个方法:

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的实现中,你可以自行决定如何把ValueKey写入目标系统,并利用PartitionIdRecordSequence等信息实现不同的处理保证(例如精确一次语义)。

此外,你必须负责 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 类型的记录(schemaTypeAVROJSONPROTOBUF_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 运行时协作有NARUber 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

可见实际内置连接器还会补充sourceConfigClasssinkConfigClass字段,用于声明连接器的配置类,提交连接器时即可通过 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-corepulsar-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 均继承该方法,因此在openwrite/read的任何位置都能上报自定义指标。指标会被 Function Worker 自动汇总并暴露给 Prometheus 抓取,配合 Grafana 即可构建连接器专属的监控面板。

总结

开发 Pulsar 连接器的完整流程可以概括为:

  1. 实现接口:Source 实现open+read,Sink 实现open+write,核心是正确封装 Record(Source 侧)并在write中妥善调用ack/fail(Sink 侧);
  2. 处理 Schema:按需选择强类型泛型、byte[]+Schema.AUTO_PRODUCE_BYTESKVRecord(键值场景)或GenericObject(Sink 通用场景);
  3. 分层测试:先单元测试,再用 testcontainers 编写端到端集成测试(参考tests/integration);
  4. 打包分发:优先 NAR(含pulsar-io.yaml元数据),或选择 Uber JAR,并处理好许可证声明;
  5. 上线监控:利用 Pulsar 指标与recordMetric自定义指标,经 Prometheus/Grafana 保障连接器健康。

从仓库源码可以看到,Pulsar 内置连接器(Kafka、RabbitMQ、Kinesis 等)全部遵循同一套接口与打包规范,因此掌握本指南的 API 契约与工程流程后,你既可以复刻内置连接器的成熟模式,也可以自由开发适配任意外部系统的自定义连接器。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

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

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

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

立即咨询