消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Pulsar IO 是 Apache Pulsar 内置的流数据处理框架它把 Pulsar 与外部系统数据库、消息队列、对象存储、数据湖等之间的数据搬运动作抽象为两类连接器ConnectorSource导入外部系统 → Pulsar与Sink导出Pulsar → 外部系统。本指南基于 Pulsar 仓库中的官方开发文档io-develop.md结合源码级实现细节完整讲解如何从零编写一个 Pulsar IO 连接器——包括实现Source/Sink接口、设计Record元数据以获得精确一次exactly-once语义、编写单元测试与集成测试以及通过 NAR 包或 Uber JAR 将连接器打包提交到 Pulsar Functions 集群。读完本文你将具备独立开发、测试与分发 Pulsar IO 连接器的完整实战能力。Pulsar IO 连接器的本质特殊化的 Pulsar FunctionsPulsar IO 连接器本质上是特殊化的 Pulsar Functions。这意味着编写一个连接器其难度与编写一个 Pulsar Function 相当你不需要关心消息如何在 Broker 与 Function 运行时之间传输也不需要实现消费与生产逻辑运行时已经替你完成了这些工作。连接器只有两种形态形态方向核心接口典型示例Source从外部系统导入数据到 PulsarSourceRabbitmqSource 将 RabbitMQ 队列的消息导入 Pulsar topicSink将 Pulsar 数据导出到外部系统SinkKinesisSink 将 Pulsar topic 的消息导出到 Kinesis 流在源码中这两个接口都位于pulsar-io/core模块的 org.apache.pulsar.io.core 包下并被标注为InterfaceAudience.Public与InterfaceStability.Stable即它们属于对外公开且稳定的公共 API开发者可以放心基于它们编写连接器。此外SourceT与SinkT都继承了AutoCloseable连接器必须正确实现close()以释放初始化时申请的资源。开发一个 Source 连接器Source 连接器的任务是从外部系统读取数据并交给 Pulsar IO 运行时写入 Pulsar。开发一个 Source核心就是实现 Source 接口。接口全貌package org.apache.pulsar.io.core; public interface SourceT extends AutoCloseable { /** * Open connector with configuration. * * param config 初始化配置 * param sourceContext source 运行所处的环境 * throws Exception IO 类型的异常 */ void open(MapString, Object config, SourceContext sourceContext) throws Exception; /** * 从 source 读取下一条消息。 * 如果没有新消息此调用应当阻塞。 * * return source 的下一条消息返回值绝不能为 null * throws Exception */ RecordT read() throws Exception; }第一步实现open方法完成初始化open在连接器初始化时只会被调用一次。在这个方法中通过传入的config参数一个MapString, Object解析连接器专属配置初始化所有必要的资源。例如Kafka 连接器可以在open中创建 Kafka 客户端。config的来源是提交连接器时通过 CLI 或 REST API传入的配置运行时会把键值对原样传入。除config之外Pulsar 运行时还会向连接器提供一个 SourceContext用于访问运行时资源例如getSourceName()获取当前 Source 的名称getOutputTopic()获取 Source 的输出 topicnewOutputMessage(topicName, schema)直接向指定 topic 发送消息返回TypedMessageBuilderTnewConsumerBuilder(schema)创建 Pulsar Consumer 构造器。SourceContext还继承了 Pulsar Functions 的 BaseContext因此可以用于记录日志、收集 metrics、读写状态等。实现者通常会把SourceContext保存下来供后续使用。第二步实现read方法读取数据read是 Source 实现者的核心任务。它的语义非常明确从外部系统读取下一条消息如果暂时没有新消息应当阻塞而不是返回空数据返回值永远不能为null没有数据就阻塞等待而不是返回 null。返回的 Record 封装了 Pulsar IO 运行时所需的全部信息。运行时拿到Record后会把它写入 Pulsar topic并按需触发 ack/fail 回调。Record携带的元数据信息RecordT中可携带的信息如下对照源码 Record.java 的默认实现信息必选/可选说明Topic NamegetTopicName可选若记录源自某个 Pulsar topic返回该 topic 名KeygetKey可选与记录关联的键用于路由与去重ValuegetValue必选记录的实际数据Partition IdgetPartitionId可选若记录源自分区数据源返回其分区 id。分区 id 将作为唯一标识的一部分供 Pulsar IO 运行时做消息去重实现精确一次处理Record SequencegetRecordSequence可选若记录源自顺序型数据源返回其记录序号。同样参与唯一标识与去重是实现精确一次的关键PropertiesgetProperties可选用户自定义属性Event TimegetEventTime可选记录在源侧的事件时间毫秒时间戳SchemagetSchema可选记录值的 Schema用于类型化写入Destination TopicgetDestinationTopic可选支持按消息粒度路由到指定 topicack与fail精确一次语义的基石除了上述元数据Record的实现还必须提供两个方法default void ack() {} // 通知 Pulsar IO 运行时该记录已成功处理 default void fail() {} // 通知 Pulsar IO 运行时该记录处理失败当 Source 侧的数据被成功写入 Pulsar 后运行时/连接器调用ack()外部系统的偏移量可以被安全提交例如 Kafka offset commit当处理失败时调用fail()外部系统的消费位置不推进消息可在恢复后重试。参考实现Kafka SourceKafkaAbstractSource 是一个极佳的参考范例。从源码可以看到它完整示范了open与Record的实现模式open中解析并校验配置KafkaSourceConfig.load(config)加载配置随后对topic、bootstrapServers、groupId、fetchMinBytes等关键参数做非空与取值范围校验如fetchMinBytes 0直接抛出IllegalArgumentException再组装 Kafka Consumer 所需的PropertiesSASL、SSL、反序列化器等创建KafkaConsumer并启动拉取线程。内部KafkaRecord实现Record接口通过getPartitionId()返回Integer.toString(record.partition())、getRecordSequence()返回record.offset()、getKey()返回 Kafka 消息 key、getValue()返回反序列化后的值。ack()完成一个CompletableFuture当autoCommitEnabledfalse时拉取线程会在所有ack完成后执行consumer.commitSync()——这正是文档所说的ack 驱动外部系统提交的落地实现。开发一个 Sink 连接器开发 Sink 连接器与开发 Source 一样简单实现 Sink 接口即可。接口全貌package org.apache.pulsar.io.core; public interface SinkT extends AutoCloseable { /** * Open connector with configuration. * * param config 初始化配置 * param sinkContext sink 运行所处的环境 * throws Exception IO 类型的异常 */ void open(MapString, Object config, SinkContext sinkContext) throws Exception; /** * 向 Sink 写入一条消息。 * * param record 要写入 sink 的记录 * throws Exception */ void write(RecordT record) throws Exception; }与 Source 对称open初始化所有必要资源例如创建数据库连接、Kafka Producer、外部 API 客户端等writeSink 实现者的核心方法接收一条Record决定如何把value以及可选的key写入目标系统。SinkContext 为 Sink 提供运行环境信息包括getSinkName()当前 Sink 的名称getInputTopics()所有输入 topic 的列表getSubscriptionType()为 Sink 提供数据的订阅类型seek(topic, partition, messageId)将指定 topic 分区的订阅重置到某个消息 ID用于故障恢复pause(topic, partition)/resume(topic, partition)暂停/恢复指定分区的消息拉取可用于背压控制。在write的实现中开发者可以决定如何把value与可选的key写入目标系统利用Record中的Partition Id、Record Sequence等信息实现不同的处理保证如借助分区/序号实现幂等写入负责消息的确认与失败成功写入后调用record.ack()写入失败时调用record.fail()让 Pulsar 侧的消费位置随之推进或回退。测试连接器连接器的测试颇具挑战性因为 Pulsar IO 连接器同时与两个系统交互Pulsar 本身以及连接器对接的外部系统二者都不容易在单测中模拟。单元测试Mock 外部服务只测连接器逻辑推荐的策略是编写高度聚焦的单元测试专注于验证连接器类自身的功能同时 Mock 掉外部服务。例如使用 Mockito 模拟外部系统的客户端验证open方法能否正确解析配置并完成初始化非法配置应抛异常read是否正确构造Recordkey、value、partition id、record sequence 是否正确write是否正确调用外部系统 API并在成功/失败时调用ack()/fail()。仓库中每个连接器模块都自带单元测试例如 TwitterFireHoseConfigTestspulsar-io/core模块下也有 SourceTest 与 SinkTest 可作为测试风格参考。集成测试testcontainers 验证端到端在单元测试覆盖充分之后强烈建议再补充独立的集成测试验证连接器与真实系统之间的端到端功能。Pulsar 项目本身使用 testcontainers 是编写连接器集成测试的绝佳范例PulsarIOTestBase.java 与 PulsarIOTestRunner.java 提供了在 Docker 容器中启动 Pulsar 集群、提交连接器并验证数据流通的基础框架sources/与sinks/子目录中分别存放着各内置 Source/Sink 的端到端测试器例如 RabbitMQSourceTester、RabbitMQSinkTester它们用真实的 RabbitMQ 容器验证消息从队列流入 Pulsar topic、再被 Sink 写回队列的完整链路。打包连接器开发和测试完成后必须将连接器打包才能提交到 Pulsar Functions 集群运行。Pulsar 支持两种打包方式。许可证义务提醒如果你打算把连接器打包分发给他人使用你有义务为自己的代码正确声明许可与版权并遵守你所用、所打包的所有第三方库的许可与版权要求。如果采用下文第一种NAR 打包方式NAR 插件会在生成的 NAR 包内自动生成DEPENDENCIES文件记录连接器所用全部库的许可与版权信息。方式一创建 NAR 包推荐NARNiFi Archive是 Apache NiFi 使用的一种自定义打包机制其核心价值是提供Java ClassLoader 隔离——不同连接器依赖的同名类库互不干扰。Pulsar 采用同一机制来打包所有的内置连接器例如 pulsar-io/aerospike/pom.xml 等内置连接器的 pom 中都引入了nifi-nar-maven-plugin。创建 NAR 包很简单只需在连接器的 Maven 项目中加入nifi-nar-maven-plugin插件plugins plugin groupIdorg.apache.nifi/groupId artifactIdnifi-nar-maven-plugin/artifactId version1.2.0/version /plugin /plugins运行mvn package后会在target/目录生成.nar文件提交连接器时直接指定该文件即可。仓库中的 Twitter FireHose 连接器是 NAR 打包的良好范例。从 TwitterFireHose.java 的源码可以看到它通过Connector(name twitter, type IOType.SOURCE, configClass TwitterFireHoseConfig.class)注解声明了连接器的元信息名称、类型、配置类并且继承的是PushSourceTweetData——一种由运行时回调驱动的推式 Source 变体。这也提醒我们pulsar-io/core中还提供了 PushSource、BatchSource、BatchPushSource 等高级基类开发者可以根据数据源的特性拉取式/推式/批式选择合适的基类从而只实现自己关心的回调方法。方式二创建 Uber JAR另一种方式是创建Uber JAR胖 JAR把连接器的全部 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两种方式都能被 Pulsar Functions 的运行时识别并加载。选择建议若面向他人分发优先使用 NAR 包——它既提供 ClassLoader 隔离避免依赖冲突又自动生成依赖许可清单Uber JAR 则适合快速迭代、自用场景无需额外插件机制但要注意依赖冲突与许可合规需要自行管理。总结开发 Pulsar IO 连接器的完整流程选型确定连接器类型Source 导入 / Sink 导出必要时从PushSource、BatchSource等基类中选择合适的父类实现接口实现open(config, context)完成配置解析与资源初始化实现read()Source或write(record)Sink完成核心数据搬运动作构造 Record正确填充 value、key、partition id、record sequence 等元数据并实现ack()/fail()以获得精确一次处理能力测试先 Mock 外部服务写单元测试再用 testcontainers 写端到端集成测试打包用nifi-nar-maven-plugin生成 NAR 包推荐或用maven-shade-plugin生成 Uber JAR提交运行将打包产物提交到 Pulsar Functions 集群相关内置连接器列表见 io-connectors。Pulsar IO 的连接器框架把两个系统之间的数据搬运收敛为实现两个 Java 接口 一个 Record让开发者能专注于业务逻辑本身而分区 id、记录序号与 ack/fail 回调的组合则为跨系统数据同步提供了可靠的消息去重与精确一次处理的基础保障。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar IO 连接器开发指南Source/Sink 接口、Schema 处理、NAR 打包与监控Apache Pulsar IO 连接器开发指南Source/Sink 接口、Schema 处理、NAR 打包与监控 本文基于 Apache Pulsar 官消息队列后端流处理Apache Pulsar Connector 开发实战指南从 Source/Sink 接口到 NAR 打包与监控Apache Pulsar Connector 开发实战指南从 Source/Sink 接口到 NAR 打包与监控 Apache Pulsar IO 框架提供消息队列后端流处理Apache Pulsar 连接器Connector开发完全指南从 Source/Sink 接口实现到 NAR 打包与监控Apache Pulsar 连接器Connector开发完全指南从 Source/Sink 接口实现到 NAR 打包与监控 本指南以 Apache Pul消息队列后端流处理上一篇libassert 开源项目教程下一篇子网路由器与出口节点Headscale Routes配置完全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考