尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Pulsar 连接器(Connector)开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控

发布时间:2026/9/24 5:25:25

资讯中心
01
ARTICLE

Apache Pulsar 连接器(Connector)开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控

Apache Pulsar 连接器(Connector)开发完全指南:从 Source/Sink 接口实现到 NAR 打包与监控
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本指南以 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将数据从外部系统导入 PulsarRabbitMQ source connector 把 RabbitMQ 队列中的消息导入 Pulsar topicSink将数据从 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 SourceT extends AutoCloseable { void open(MapString, Object config, SourceContext sourceContext) throws Exception; RecordT 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 MapString, Object config, SourceContext sourceContext) throws Exception;open方法在 Source 连接器初始化时被调用。在这个方法里通过传入的configMap读取该连接器所有的自定义配置项并初始化所需资源如网络客户端、连接池等将SourceContext保存下来供后续使用。例如 Kafka Source 就在open中创建了 Kafka 消费客户端。参考 KafkaAbstractSource.java其open方法内部做了三件事用KafkaSourceConfig.load(config)从config反序列化出强类型的配置对象对topic、bootstrapServers、groupId等关键配置做Objects.requireNonNull非空校验并对fetchMinBytes、sessionTimeoutMs、heartbeatIntervalMs等参数做合法性校验将配置拼装成 KafkaProperties创建KafkaConsumer并启动拉取线程start()。这种config 解析 → 参数校验 → 资源初始化的三段式写法是内置连接器普遍遵循的模式。除了获取配置Pulsar 运行时还通过 SourceContext 向连接器暴露运行环境信息例如getSourceName()/getOutputTopic()当前 Source 的名称与输出 topicnewOutputMessage(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 */ RecordT read() throws Exception;read方法负责从外部系统读取下一条消息。当没有新消息时实现应当阻塞等待而不是返回null——这是该方法最重要的约定。返回的 Record 需要封装 Pulsar IO 运行时所需的如下信息Record 变量变量必填说明TopicName否该记录来源的 Pulsar topic 名Key否消息可选的键。键用于路由详见 Routing modesValue是记录的实际数据EventTime否记录在源端的事件时间自 epoch 起的毫秒数PartitionId否若记录来自分区源返回其分区 ID。Pulsar IO 运行时用它与RecordSequence一起构成唯一标识用于去重并实现 exactly-once 处理保证RecordSequence否若记录来自有序源返回其序号同样是去重/精确一次语义的组成部分Properties否用户自定义属性MapString, StringDestinationTopic否该消息应被写入的目标 topic支持按消息粒度路由Message否携带用户发送数据的Message对象见 Message.javaRecord 方法方法说明ack确认该记录已被完全处理fail标记该记录处理失败Kafka Source 的内部类 KafkaRecord 给出了很好的参考实现getPartitionId()/getPartitionIndex()返回 Kafka 分区号getRecordSequence()返回 Kafka offsetgetKey()返回消息键getValue()返回反序列化后的值ack()通过CompletableFuture完成回调驱动 Kafka 消费位移提交。处理 Schema 信息Source 侧Pulsar IO 会自动处理 Schema并基于 Java 泛型提供强类型 API。如果你明确知道自己要产出的数据类型可以在 Source 声明中直接指定对应的 Java 类型public class MySource implements SourceString { public RecordString read() {} }如果你要实现的 Source 需要兼容任意 Schema可以改用byte[]或ByteBuffer配合Schema.AUTO_PRODUCE_BYTES()public class MySource implements Sourcebyte[] { public Recordbyte[] read() { Schema wantedSchema .... Recordbyte[] myRecord new MyRecordImplementation(); .... } class MyRecordImplementation implements Recordbyte[] { public byte[] getValue() { return ....encoded byte[]...that represents the value } public Schemabyte[] getSchema() { return Schema.AUTO_PRODUCE_BYTES(wantedSchema); } } }正确处理 KeyValue 类型要正确处理KeyValue类型你的 Record 实现需要遵循三条规则实现 KVRecord 接口并实现getKeySchema()、getValueSchema()和getKeyValueEncodingType()三个方法Record.getValue()必须返回KeyValue对象Record.getSchema()可以返回null。当 Pulsar IO 运行时遇到KVRecord时会自动完成以下转换正确设置KeyValueSchema按照KeyValueEncodingSEPARATED或INLINE编码消息键与消息值。KVRecordK, V接口本身位于pulsar-functions/api-java其getKeyValueEncodingType()返回org.apache.pulsar.common.schema.KeyValueEncodingType枚举。Kafka Source 的 KeyValueKafkaRecord 就是典型实现它实现了KVRecordObject, Object分别持有 keySchema 与 valueSchema并返回KeyValueEncodingType.SEPARATED键值分离编码。关于如何实现一个完整的 Source 连接器可以参考内置的 KafkaSource。开发 Sink 连接器开发 Sink 连接器与开发 Source 连接器非常类似实现 Sink 接口即实现open与write两个方法public interface SinkT extends AutoCloseable { void open(MapString, Object config, SinkContext sinkContext) throws Exception; void write(RecordT 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 MapString, Object config, SinkContext sinkContext) throws Exception;与 Source 的open类似这里读取配置、初始化外部系统客户端等资源。SinkContext提供 Sink 运行环境包括getSinkName()Sink 名称getInputTopics()所有输入 topic 列表getSubscriptionType()订阅类型seek(topic, partition, messageId)将订阅重置到指定消息 IDpause(topic, partition)/resume(topic, partition)暂停/恢复消费指定 topic 分区同样继承BaseContext可获得日志、密钥、状态存储、计数器、指标上报等能力。实现 write 方法/** * Write a message to Sink * param record record to write to sink * throws Exception */ void write(RecordT 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 SinkString { public void write(RecordString record) {} }若需要实现可适配任意 Schema 的 Sink可以使用特殊的GenericObject接口public class MySink implements SinkGenericObject { public void write(RecordGenericObject 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 与值 Schemapublic class MySink implements SinkGenericObject { public void write(RecordGenericObject 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推荐NARNiFi Archive是 Apache NiFi 使用的一种自定义打包机制用于提供一定程度的 Java ClassLoader 隔离。Pulsar 用同样的机制打包全部内置连接器见 pulsar-io 下各连接器模块。打包连接器最简单的方式是使用 nifi-nar-maven-plugin 创建 NAR 包。在连接器模块的 Maven 工程中加入该插件plugins plugin groupIdorg.apache.nifi/groupId artifactIdnifi-nar-maven-plugin/artifactId version1.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.yamlname: 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——该模块在buildplugins中仅声明nifi-nar-maven-plugin并依赖pulsar-io-core与pulsar-io-common等模块即可产出可被 Pulsar 加载的 NAR 包。方式二Uber JAR另一种做法是创建包含连接器全部 JAR 文件及其他资源文件的uber JAR无需任何目录内部结构。使用 maven-shade-plugin 构建plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.1.1/version executions execution phasepackage/phase goals goalshade/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 SinkString { Override public void open(MapString, Object config, SinkContext sinkContext) throws Exception { sinkContext.recordMetric(foo, 1); } Override public void write(RecordString 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 实现openreadSink 实现openwrite核心是正确封装 RecordSource 侧并在write中妥善调用ack/failSink 侧处理 Schema按需选择强类型泛型、byte[]Schema.AUTO_PRODUCE_BYTES、KVRecord键值场景或GenericObjectSink 通用场景分层测试先单元测试再用 testcontainers 编写端到端集成测试参考tests/integration打包分发优先 NAR含pulsar-io.yaml元数据或选择 Uber JAR并处理好许可证声明上线监控利用 Pulsar 指标与recordMetric自定义指标经 Prometheus/Grafana 保障连接器健康。从仓库源码可以看到Pulsar 内置连接器Kafka、RabbitMQ、Kinesis 等全部遵循同一套接口与打包规范因此掌握本指南的 API 契约与工程流程后你既可以复刻内置连接器的成熟模式也可以自由开发适配任意外部系统的自定义连接器。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。