
Apache Pulsar IO 连接器全解析Source、Sink 与处理保证机制实战指南【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar消息系统只有在能够轻松与数据库、其他消息系统等外部系统对接时才能发挥最大价值。Pulsar IOPulsar 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 将外部数据流入 PulsarSink 将 Pulsar 数据流出到外部系统)核心概念Source 与 SinkPulsar IO 连接器分为两种类型source源和sink汇二者共同构成外部系统 ⇄ Pulsar的双向数据通道。Source将外部数据流入 PulsarSources将外部系统的数据注入 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.protoenum 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。在配置模型中该字段存在于SourceConfigSourceConfig.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.mdsinks子命令管理 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 负责将外部数据流入 PulsarSink 负责将 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),仅供参考