Apache Pulsar IO 连接器开发指南:Source/Sink 接口、Schema 处理、NAR 打包与监控
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文基于 Apache Pulsar 官方文档《How to develop Pulsar connectors》系统讲解如何开发 Pulsar 连接器IO Connector实现Source/Sink接口、Record数据模型的完整字段、Schema含 KeyValue处理规则、测试策略以及 NAR 与 Uber JAR 两种打包方式和运行时监控方法。读完本文你可以独立编写一个可从外部系统向 Pulsar 导入数据Source或从 Pulsar 导出数据Sink的连接器完成打包部署并通过指标进行健康监测。Pulsar 连接器是什么Pulsar 连接器用于在 Pulsar 与其他系统之间搬运数据。由于连接器本质上是运行在 Pulsar Functions 运行时上的特殊 Function因此开发连接器的流程与开发 Pulsar Function 类似。连接器分为两种类型类型说明典型示例Source从外部系统向 Pulsar 导入数据RabbitMQ Source 连接器将 RabbitMQ 队列中的消息导入 Pulsar topicSink从 Pulsar 向外部系统导出数据Kinesis Sink 连接器将 Pulsar topic 中的消息导出到 Kinesis 数据流核心接口定义在pulsar-io模块中Source 接口pulsar-io/core/src/main/java/org/apache/pulsar/io/core/Source.javaSink 接口pulsar-io/core/src/main/java/org/apache/pulsar/io/core/Sink.java两个接口均标注InterfaceAudience.Public和InterfaceStability.Stable且都继承AutoCloseable即实现方还需实现close()方法释放资源。开发 Source 连接器开发 Source 连接器就是实现SourceT接口核心是实现open和read两个方法。1. 实现 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;该方法在 Source 连接器初始化时被调用。你可以在这里通过传入的config参数MapString, Object获取所有连接器专属配置并初始化所需资源。例如 Kafka 连接器就可以在open方法中创建 Kafka client保存SourceContext以便后续使用。从源码看SourceContext见 SourceContext.java继承自BaseContext并提供以下能力方法作用getSourceName()获取当前 Source 的名称getOutputTopic()获取 Source 的输出 topic 名称newOutputMessage(topicName, schema)基于指定 Schema 构造发往指定 topic 的消息 buildernewConsumerBuilder(schema)创建带 Schema 的 ConsumerBuilder这意味着 Source 可以直接把读取到的数据通过newOutputMessage写到指定 topic而不必依赖默认的单一输出 topic——这正对应下表中DestinationTopic字段的能力。2. 实现 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;关键约定当没有数据可返回时实现应当阻塞而不是返回null。返回的 Record 需要封装 Pulsar IO 运行时所需的信息。Record应提供以下变量getter变量必填说明TopicName否该记录来源的 Pulsar topic 名称Key否消息可选的键可参见消息路由模式Value是记录的实际数据EventTime否来源系统中的事件时间PartitionId否若记录来自分区式源则返回其PartitionId。它被 Pulsar IO 运行时用作唯一标识符的一部分用于消息去重并实现 exactly-once 处理保证RecordSequence否若记录来自顺序源则返回其RecordSequence。同样被运行时用作唯一标识符的一部分用于去重并实现 exactly-once 保证Properties否记录携带的用户自定义属性DestinationTopic否该消息应当写入的 topic支持逐消息路由Message否携带用户发送数据的类Pulsar client 的Message见 Message.javaRecord还应提供以下方法方法说明ack确认该记录已被完全处理fail表示该记录处理失败对照 Record.java 源码可以看到除getValue()外上述字段和方法均提供了默认实现返回Optional.empty()或空集合因此实现方只需覆盖自己关心的部分。参考实现Kafka Source官方提供的完整 Source 实现参考是 KafkaAbstractSource其实现方式很好地印证了上述约定open()约 L61-L127先用KafkaSourceConfig.load(config)解析配置并校验topic、bootstrapServers、groupId等必填项未设置直接抛异常随后构造Properties创建KafkaConsumer从源码结构看Kafka Source 采用了推push模型start()约 L149-L183中启动一个名为 Kafka Source Thread 的 runner 线程循环consumer.poll(Duration.ofSeconds(1L))拉取ConsumerRecords逐条构造KafkaRecord后推送给上层read()则由父类KafkaPushSource的内部队列承接。这正是应对read 必须阻塞约束的常见做法——外部系统若是推送式消息源Kafka、Twitter Firehose用一个后台线程轮询 队列桥接即可内部类KafkaRecord约 L188-L235展示了如何填充去重所需的元数据getPartitionId()返回 Kafka 分区号字符串、getPartitionIndex()返回分区序号、getRecordSequence()返回 offsetack()则完成内部CompletableFuture由轮询线程在CompletableFuture.allOf(futures).get()之后执行consumer.commitSync()未开启自动提交时从而把Pulsar 侧写成功与Kafka 侧提交 offset关联起来。如果你的源支持 offset 语义务必像这样把外部系统的分区号与 offset 映射到PartitionId/RecordSequence这是获得 exactly-once 保证的前提。Source 的 Schema 处理Pulsar IO 会自动处理 Schema并基于 Java 泛型提供强类型 API。如果你已知要生产的 Schema 类型直接在 Source 声明中用对应的 Java 类即可public class MySource implements SourceString { public RecordString read() {} }如果你想实现一个兼容任意 Schema 的 Source可以用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类型记录实现需要遵循以下约定必须实现 KVRecord 接口并实现getKeySchema()、getValueSchema()和getKeyValueEncodingType()三个方法源码见 KVRecord.java必须以KeyValue对象作为Record.getValue()的返回值Record.getSchema()可以返回 null。当 Pulsar IO 运行时遇到KVRecord时会自动正确设置KeyValueSchema按KeyValueEncodingSEPARATED或INLINE对消息 Key 和消息 Value 进行编码。仓库中的KeyValueKafkaRecord位于 KafkaAbstractSource.java就是一个真实例子它继承KafkaRecord并实现KVRecordObject, ObjectgetValueSchema()/getKeySchema()分别返回 value 与 key 的 SchemagetKeyValueEncodingType()返回KeyValueEncodingType.SEPARATED。开发 Sink 连接器开发 Sink 连接器与开发 Source 连接器类似实现 Sink 接口即实现open和write方法。1. 实现 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;2. 实现 write 方法/** * Write a message to Sink * param record record to write to sink * throws Exception */ void write(RecordT record) throws Exception;实现时你可以自行决定如何将Value和Key写入外部系统并利用PartitionId、RecordSequence等所有提供的信息实现不同的处理保证例如基于 offset 的去重与顺序写。你还需要在消息发送成功后调用record.ack()在发送失败时调用record.fail()。从源码看SinkContext 除了getSinkName()和getInputTopics()获取所有输入 topic外还提供流控能力pause(topic, partition)/resume(topic, partition)可暂停与恢复指定 topic、分区的消息请求常用于下游写入过慢时背压seek(topic, partition, messageId)可将订阅重置到指定消息位点。这些 default 方法在未实现的上下文中会抛出UnsupportedOperationException实现方可按需覆盖。Sink 的 Schema 处理Pulsar IO 自动处理 Schema并提供基于 Java 泛型的强类型 API。如果你已知要消费的 Schema 类型直接在 Sink 声明中声明对应 Java 类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 记录schemaTypeAVRO、JSON、PROTOBUF_NATIVE可以把genericObject变量转换为GenericRecord使用getFields()和getField()API也可以通过genericObject.getNativeObject()访问原生的 AVRO 记录。对于 KeyValue 类型可以同时访问 key 与 value 的 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(); } ... } }测试连接器测试连接器具有挑战性因为 Pulsar IO 连接器要与两个系统交互而这两个系统Pulsar 和连接器对接的外部系统都可能难以 mock。建议在 mock 外部服务的前提下为连接器功能编写专门的测试单元测试为连接器创建常规单元测试覆盖配置解析、Record 构造、Schema 处理等纯逻辑集成测试在拥有足够单元测试后添加独立的集成测试验证端到端功能。Pulsar 的所有集成测试都基于 testcontainers容器化测试框架。仓库中的集成测试可以参考 tests/integration/src/test/java/org/apache/pulsar/tests/integration/io 目录其中包含 PulsarIOTestBase.java集成测试基类、PulsarIOTestRunner.java以及RabbitMQSourceTester.java、RabbitMQSinkTester.java等基于容器化的真实连接器测试并带有sources/、sinks/子目录分类可作为新连接器集成测试的模板。打包连接器开发并测试完连接器后需要将其打包才能提交到 Pulsar Functions 集群运行。有两种方式与 Pulsar Functions 运行时配合NAR 和 Uber JAR。注意如果你计划打包分发连接器供他人使用你有义务为自己的代码以及代码所用到的所有库和分发包正确添加 license 与 copyright。若使用 NAR 方式NAR 插件会自动在生成的 NAR 包中创建DEPENDENCIES文件包含连接器所有依赖库的 license 与版权信息。方式一NAR 打包NARNiFi Archive是一种源自 Apache NiFi 的自定义打包机制用于提供一定程度的 Java ClassLoader 隔离。Pulsar 的所有内置连接器见连接器列表都用这种方式打包。打包 Pulsar 连接器最简单的方法是在 Maven 项目中引入nifi-nar-maven-pluginplugins plugin groupIdorg.apache.nifi/groupId artifactIdnifi-nar-maven-plugin/artifactId version1.2.0/version /plugin /plugins版本说明本文档示例给出 1.2.0在当前仓库中根 pom.xml 通过属性nifi-nar-maven-plugin.version统一管理版本当前为 1.5.0各连接器模块的 pom如 pulsar-io/twitter/pom.xml直接声明插件而不写版本号由父 POM 继承管理。独立项目自行引入时可显式指定版本。同时必须创建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)仓库中的真实示例见 pulsar-io/twitter/src/main/resources/META-INF/services/pulsar-io.yamlname: twitter description: Ingest data from Twitter firehose sourceClass: org.apache.pulsar.io.twitter.TwitterFireHose sourceConfigClass: org.apache.pulsar.io.twitter.TwitterFireHoseConfig可以看到除文档描述的name、description、sourceClass外实际项目还会声明sourceConfigClassSink 侧对应sinkConfigClass运行时据此对配置做类型化解析。使用 Gradle 的用户可以使用 Gradle Plugin Portal 上的 Gradle Nar 插件gradle-nar-plugin。关于 NAR 的使用可参考 pulsar-io/twitter/pom.xml。方式二Uber JAR另一种方式是创建一个包含连接器所有 JAR 文件与其他资源文件的uber JAR无需任何内部目录结构。可以使用maven-shade-plugin创建 uber JARplugin 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其中filters部分可按需配置 includes/excludes 规则来控制打入 uber JAR 的内容*:*表示作用于所有 artifact注意不带includes/excludes的 filter 本身不产生实际过滤效果实际项目请按 shade 插件文档补充具体的包含/排除规则。监控连接器Pulsar 连接器让你轻松地在 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(name, value)是BaseContextSourceContext/SinkContext的父接口提供的能力可在open、read/write等任意阶段调用将业务计数如成功写出条数外部接口失败次数暴露为 Prometheus 指标与 Pulsar 内建指标共同构成连接器的可观测体系。小结开发一个完整的 Pulsar 连接器包含五个环节选择Source外部 → Pulsar或SinkPulsar → 外部实现openread/write用Record封装数据必要时填充PartitionId/RecordSequence/Key以支持 exactly-once 与路由用 Java 泛型声明 Schema或用byte[]/GenericObjectKVRecord编写通用型连接器编写 mock 外部服务的单元测试与基于 testcontainers 的集成测试用 NAR推荐自动处理 ClassLoader 隔离与依赖许可声明或 Uber JAR 打包配合pulsar-io.yaml元数据部署并通过内建指标与recordMetric自定义指标监控运行状态。关键入口文件Source.java、Sink.java、Record.java、KafkaAbstractSource.java、集成测试目录。赞分享消息队列后端流处理【免费下载链接】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 Pul消息队列后端流处理架构级实战指南3步实现Redis系统零停机升级的高可用方案架构级实战指南3步实现Redis系统零停机升级的高可用方案 在现代微服务架构中 系统高可用 已成为企业级应用的基石而 零停机部署 和 服务平滑升级 则是保消息队列后端流处理Apache Pulsar Connector 开发完全指南从 Source/Sink 接口实现到 NAR 打包与监控Apache Pulsar Connector 开发完全指南从 Source/Sink 接口实现到 NAR 打包与监控 本指南以 Apache Pulsar消息队列后端流处理上一篇Rolldown Module ID 解析字符串路径身份的归一化设计与跨平台一致性下一篇Zotero Duplicates Merger学术文献库智能去重解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考