news 2026/9/23 10:13:07

Apache Pulsar IO 连接器全解析:Source、Sink 与处理保证机制实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar IO 连接器全解析:Source、Sink 与处理保证机制实战指南

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-onceat-least-onceeffectively-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_ONCE
  • ATMOST_ONCE
  • EFFECTIVELY_ONCE

若创建连接器时未指定--processing-guarantees,默认语义为ATLEAST_ONCE

在配置模型中,该字段存在于SourceConfig(SourceConfig.java)与对应的SinkConfig中,类型即FunctionConfig.ProcessingGuarantees,创建与更新连接器时均可通过 Admin CLI 传入。

创建连接器时设置

以下以Admin CLI为例。关于REST APIJava 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_ONCEATMOST_ONCEEFFECTIVELY_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 完成,其下分sourcessinks两个子命令组:

  • 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枚举区分FUNCTIONSOURCESINK),由同一套运行时框架调度执行,只是组件类型与数据处理方向不同。

关于 Functions worker 的独立运行方式,参见 Functions worker。

总结

Pulsar IO 通过 Source 与 Sink 两类连接器打通了 Pulsar 与外部系统之间的数据通路:Source 负责将外部数据流入 Pulsar,Sink 负责将 Pulsar 数据流出到外部系统。在处理保证方面,at-most-onceat-least-onceeffectively-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 10:07:37

阿里云ROS Terraform托管服务:企业级IaC运行时底座

1. ROS Terraform 托管服务不是“ROS机器人操作系统”的缩写,而是阿里云资源编排服务(Resource Orchestration Service)的官方命名刚看到标题里“ROS Terraform”这个组合,很多刚接触云基础设施的同学第一反应是:“ROS…

作者头像 李华
网站建设 2026/9/23 10:06:37

存储过程——游标:从 OPEN-FOR 到 FOR 循环的 TaoToken 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/23 10:05:47

Python代码质量检查工具Flake8详解与应用指南

1. 为什么我们需要代码检查工具在编写Python代码时,即使是最有经验的开发者也会不经意间引入各种问题。从简单的空格错误到潜在的逻辑缺陷,这些"小问题"往往会像滚雪球一样,最终导致难以调试的错误。这就是为什么我们需要像Flake8这…

作者头像 李华
网站建设 2026/9/23 10:04:49

UI生成新路径:本地化专用小模型微调与部署实战

最近半年,群里聊AI做UI的频率明显降下来了,不是说不做了,而是没人再拿着一张大模型通用对话去生成整页界面了。以前大家喜欢把需求一长串丢给在线大模型,让它直接“写一个后台页面”;现在更多强调的是:单独…

作者头像 李华
网站建设 2026/9/23 10:04:07

2026年配音工具技术选型:7款实测,从免费试听到批量生成全链路

配音软件哪个好用?做技术教程或批量内容生产时,这个问题几乎每个月都会被问一遍。自己录环境不允许,外包成本高,AI配音工具又参差不齐。2026年,TTS市场已经分层清晰:轻量免费工具满足个人创作者快速出稿&am…

作者头像 李华