Apache Pulsar IO 连接器全解析:Source、Sink 与处理保证机制实战指南
【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar
消息系统只有在能够轻松与数据库、其他消息系统等外部系统对接时,才能发挥最大价值。Pulsar IO(Pulsar Connector)正是承担这一职责的官方连接器框架,它让开发者可以创建、部署和管理与外部系统交互的连接器,例如 Apache Cassandra、Aerospike 等。
本文以 io-overview.md 为核心主线,系统讲解 Pulsar IO 的两大类连接器(Source 与 Sink)的架构概念、三种处理保证(Processing Guarantees)的语义与配置方式,并结合本仓库源码(proto 定义、Source/Sink 配置模型与运行时代码)深入剖析其底层实现原理,最后给出通过 Connector Admin CLI 管理连接器的完整路径。读完本文,你将能够:理解 Source/Sink 的职责边界与数据流向,准确选择并配置at-most-once、at-least-once、effectively-once三种语义,以及使用 Admin CLI 完成连接器的创建、更新、启停等全生命周期管理。
Pulsar IO 架构示意图:Source 将外部数据流入 Pulsar,Sink 将 Pulsar 数据流出到外部系统")
核心概念:Source 与 Sink
Pulsar IO 连接器分为两种类型:source(源)和sink(汇),二者共同构成"外部系统 ⇄ Pulsar"的双向数据通道。
Source:将外部数据流入 Pulsar
Sources将外部系统的数据注入 Pulsar。
Source 扮演"生产者"角色,负责从外部系统拉取数据并写入 Pulsar topic。常见的外部来源包括:
- 其他消息系统(如 Kafka、RabbitMQ 等);
- 日志与监控数据管道(firehose 风格的数据管道 API,如 Kinesis、Twitter 等)。
例如仓库 pulsar-io/kafka 提供了从 Kafka 消费数据并写入 Pulsar 的 Source 实现,pulsar-io/kinesis 则对接 AWS Kinesis。
关于 Pulsar 内置 Source 连接器的完整清单,参见 source connector。
Sink:将 Pulsar 数据流出到外部系统
Sinks将 Pulsar 的数据送入外部系统。
Sink 扮演"消费者"角色,订阅 Pulsar topic,并把消息写入外部系统。常见的外部目标包括:
- 其他消息系统;
- SQL 与 NoSQL 数据库(Cassandra、Aerospike、ElasticSearch、JDBC 兼容数据库等)。
仓库 pulsar-io/cassandra、pulsar-io/aerospike、pulsar-io/elastic-search、pulsar-io/jdbc 等即是对应 Sink 的实现。
关于 Pulsar 内置 Sink 连接器的完整清单,参见 sink connector。
处理保证(Processing Guarantees)
处理保证(Processing Guarantees)用于定义向 Pulsar topic 写入消息时发生错误如何处理。需要特别强调的是:Pulsar 连接器(Connectors)与函数(Functions)使用完全相同的处理保证机制——因为从实现层面看,Source、Sink 与 Function 都是运行在 Functions worker 上的实例组件(下文详述)。
三种投递语义
| 投递语义 | 描述 |
|---|---|
at-most-once | 发送给连接器的每条消息最多被处理一次(可能被处理一次,也可能完全不被处理) |
at-least-once | 发送给连接器的每条消息至少被处理一次(可能被处理一次,也可能被处理多次) |
effectively-once | 发送给连接器的每条消息只对应一个输出(即最终结果恰好一次) |
保证的边界:Pulsar 与外部系统各负其责
处理保证并非只依赖 Pulsar 单方面就能实现,它同时取决于外部系统以及 Source/Sink 的具体实现:
- Source 侧:Pulsar 保证"向 Pulsar topic 写入消息"遵守所选的处理保证。写入行为完全在 Pulsar 控制范围内,因此这一侧语义由 Pulsar 自身负责兑现。
- Sink 侧:处理保证取决于 Sink 的实现。如果 Sink 的实现没有以幂等方式处理重试,那么该 Sink 就无法兑现所声明的处理保证。例如在
at-least-once语义下,消息可能被重复投递给外部系统,Sink 必须自行实现去重或幂等写入才能保证最终一致性。
底层实现:源码中的枚举定义
从源码层面看,三种语义被建模为 proto 枚举,定义在 Function.proto:
enum ProcessingGuarantees { ATLEAST_ONCE = 0; // [default value] ATMOST_ONCE = 1; EFFECTIVELY_ONCE = 2; }注意ATLEAST_ONCE = 0被显式标注为默认值——这与文档中"未指定时默认ATLEAST_ONCE"的约定完全一致。该枚举同时作为FunctionDetails消息的字段(ProcessingGuarantees processingGuarantees = 6;,见 Function.proto),说明连接器与函数在处理保证上共用同一套运行时契约。
实现细节:Source 侧如何兑现语义
从运行时源码可以直观看到 Pulsar 如何在 Source 侧兑现不同语义。在 PulsarSource.java 中,每条消息被包装为PulsarRecord,其 ack 与 fail 回调行为随处理保证不同而不同:
- 当语义为
EFFECTIVELY_ONCE时:- 成功处理 → 调用
consumer.acknowledgeCumulativeAsync(message)(累积确认,一次性确认当前及之前的所有消息); - 处理失败 → 直接抛出
RuntimeException,不做负确认,等待框架层面的兜底,从而避免消息被重复投递;
- 成功处理 → 调用
- 其他语义(
ATLEAST_ONCE/ATMOST_ONCE):- 成功处理 → 调用
consumer.acknowledgeAsync(message)(单条确认); - 处理失败 → 调用
consumer.negativeAcknowledge(message)(负确认),使消息后续被重新投递,从而支撑at-least-once语义。
- 成功处理 → 调用
这段实现也印证了原文档的结论:Source 侧的保证是 Pulsar 能够完全控制的,因为确认(ack)与负确认(nack)的时机完全由 Pulsar 客户端代码决定。
设置(Set)与更新(Update)处理保证
可选的语义值
创建或更新连接器时,可以使用以下三种语义值(与 proto 枚举一一对应):
ATLEAST_ONCEATMOST_ONCEEFFECTIVELY_ONCE
若创建连接器时未指定
--processing-guarantees,默认语义为ATLEAST_ONCE。
在配置模型中,该字段存在于SourceConfig(SourceConfig.java)与对应的SinkConfig中,类型即FunctionConfig.ProcessingGuarantees,创建与更新连接器时均可通过 Admin CLI 传入。
创建连接器时设置
以下以Admin CLI为例。关于REST API与Java Admin API的对应用法,参见 io-use.md。
Source 示例(创建时指定ATMOST_ONCE):
$ bin/pulsar-admin sources create \ --processing-guarantees ATMOST_ONCE \ # Other source configs关于pulsar-admin sources create的更多选项,参见 reference-connector-admin.md。
Sink 示例(创建时指定EFFECTIVELY_ONCE):
$ bin/pulsar-admin sinks create \ --processing-guarantees EFFECTIVELY_ONCE \ # Other sink configs关于pulsar-admin sinks create的更多选项,参见 reference-connector-admin.md。
更新连接器的处理保证
连接器创建之后,仍可随时更新其处理保证,更新的语义值同样是ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE三者之一。
Source 示例(将处理保证更新为EFFECTIVELY_ONCE):
$ bin/pulsar-admin sources update \ --processing-guarantees EFFECTIVELY_ONCE \ # Other source configs关于pulsar-admin sources update的更多选项,参见 reference-connector-admin.md。
Sink 示例(将处理保证更新为ATMOST_ONCE):
$ bin/pulsar-admin sinks update \ --processing-guarantees ATMOST_ONCE \ # Other sink configs关于pulsar-admin sinks update的更多选项,参见 reference-connector-admin.md。
管理与运行连接器
通过 Connector Admin CLI 管理
连接器的全生命周期管理——创建(create)、更新(update)、启动(start)、停止(stop)、重启(restart)、重新加载(reload)、删除(delete)等操作——统一通过 Connector Admin CLI 完成,其下分sources与sinks两个子命令组:
sources子命令:管理 Source 连接器,参见 io-cli.md;sinks子命令:管理 Sink 连接器,参见 io-cli.md。
运行机制:连接器与 Functions 共享 Functions worker
连接器(Source 与 Sink)和函数(Functions)都是实例(instance)的组成部分,全部运行在 Functions worker 之上。当你通过 Connector Admin CLI 或 Functions Admin CLI 管理某个 Source、Sink 或 Function 时,实际上是在某个 worker 上启动一个实例。
这也解释了原文档开头"连接器与 Functions 使用相同的处理保证"的深层原因:Source、Sink 与 Function 在运行时统一被描述为FunctionDetails(见 Function.proto,其中ComponentType枚举区分FUNCTION、SOURCE、SINK),由同一套运行时框架调度执行,只是组件类型与数据处理方向不同。
关于 Functions worker 的独立运行方式,参见 Functions worker。
总结
Pulsar IO 通过 Source 与 Sink 两类连接器打通了 Pulsar 与外部系统之间的数据通路:Source 负责将外部数据流入 Pulsar,Sink 负责将 Pulsar 数据流出到外部系统。在处理保证方面,at-most-once、at-least-once与effectively-once三种语义由 Pulsar 与连接器实现共同兑现——Source 侧的保证完全由 Pulsar 控制(如 PulsarSource.java 中按语义选择累积确认或负确认),Sink 侧则依赖外部写入的幂等性。开发者在创建与更新连接器时,通过pulsar-admin sources/sinks create|update --processing-guarantees即可完成语义配置,并借助 Connector Admin CLI 对运行在 Functions worker 上的连接器实例进行全生命周期管理。
【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考