- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Apache Pulsar 的 IO 连接器(Pulsar IO Connectors)用于打通 Pulsar 与外部数据系统(数据库、其他消息系统等),是构建数据管道的关键能力。本文以官方 2.3.2 版文档为主体,结合本仓库源码,系统讲解 Source / Sink 概念、三种处理保证(Processing Guarantees)语义及其设置与更新方法,并延伸介绍连接器的部署、管理与监控实操。读完本文,你将能够理解 Pulsar IO 的架构模型,掌握通过pulsar-adminCLI 创建、更新、删除、查看 Source 与 Sink 的完整操作方式。
Pulsar IO 连接器:消息系统与外部世界之间的桥梁
消息系统(Messaging System)只有能轻松与数据库、其他消息系统等外部系统对接时,才最能发挥其价值。Pulsar IO 连接器正是为此而生:它让你能够轻松地创建(create)、部署(deploy)和管理(manage)与外部系统交互的连接器,例如 Apache Cassandra、Aerospike 以及许多其他系统。
从整体上看,Pulsar IO 连接器在功能形态上属于 Pulsar Functions 体系:连接器与函数(Functions)都是实例(instance)的组成部分,并且都运行在 Functions Worker 上。当你通过 Connector Admin CLI(即pulsar-admin sources/pulsar-admin sinks)或 Functions Admin CLI 管理一个 Source、Sink 或 Function 时,会在某个 Worker 上启动一个实例。关于 Functions Worker 的更多细节可参考 Functions worker。
Source 与 Sink:连接器的两种类型
Pulsar IO 连接器分为两种类型:Source(数据源)和Sink(数据汇)。下面的示意图展示了 Source、Pulsar 与 Sink 之间的关系:
Pulsar IO 连接器(Source 与 Sink)架构示意图")
Source:将外部数据流入 Pulsar
Source 的职责是将外部系统中的数据馈入 Pulsar(feed data from external systems into Pulsar)。
常见的 Source 包括其他消息系统以及 firehose 风格的数据管道 API(例如 Twitter Firehose)。一个典型的例子是 Kafka Source:它从 Kafka 主题消费消息,再写入 Pulsar 主题。
Pulsar 内置 Source 连接器的完整清单,参见 source connector。从本仓库的目录结构看,内置连接器分布在 pulsar-io 目录下,例如 Kafka Source、RabbitMQ Source、Twitter Firehose 等,每个连接器都是一个独立的 Maven 模块。
Sink:将 Pulsar 数据流出到外部系统
Sink 的职责是将 Pulsar 中的数据馈入外部系统(feed data from Pulsar into external systems)。
常见的 Sink 包括其他消息系统以及 SQL、NoSQL 数据库,例如 Cassandra Sink、ElasticSearch Sink、HBase Sink、Redis Sink、MongoDB Sink、InfluxDB Sink 等。
Pulsar 内置 Sink 连接器的完整清单,参见 sink connector。
内置连接器清单(版本 2.3.2)
根据 Builtin Connectors 文档,Pulsar 发行版内置了经过打包和测试的一组常用连接器,覆盖大多数常见数据系统:
- Aerospike Sink Connector
- Cassandra Sink Connector
- Kafka Sink Connector / Kafka Source Connector
- Kinesis Sink Connector
- RabbitMQ Source Connector / RabbitMQ Sink Connector
- Twitter Firehose Source Connector
- CDC Source Connector(基于 Debezium)
- Netty Source Connector
- HBase Sink Connector
- ElasticSearch Sink Connector
- File Source Connector
- HDFS Sink Connector
- MongoDB Sink Connector
- Redis Sink Connector
- Solr Sink Connector
- InfluxDB Sink Connector
处理保证(Processing Guarantees):连接器的消息语义
处理保证(Processing Guarantees)用于处理向 Pulsar 主题写入消息时发生的错误。需要特别注意的是:Pulsar 连接器(Source 与 Sink)与 Pulsar Functions 使用相同的处理保证语义,具体如下表所示:
| 交付语义(Delivery semantic) | 描述(Description) |
|---|---|
at-most-once(至多一次) | 发送给连接器的每条消息要么被处理一次,要么不被处理。 |
at-least-once(至少一次) | 发送给连接器的每条消息可能被处理一次,也可能被处理多次。 |
effectively-once(恰好一次 / 有效一次) | 发送给连接器的每条消息对应一个输出。 |
这三种语义同样出现在 Pulsar Functions 的处理保证文档 functions-guarantees.md 中,二者保持一致。
处理保证的适用范围与边界
需要强调的是:连接器的处理保证不仅依赖于 Pulsar 自身的保证,还与外部系统的实现(即 Source 和 Sink 的具体实现)密切相关。
- Source:Pulsar 保证向 Pulsar 主题写入消息时遵守处理保证。这一步完全在 Pulsar 的控制范围之内。
- Sink:处理保证依赖于 Sink 的实现。如果 Sink 实现没有以幂等(idempotent)的方式处理重试,那么该 Sink 就无法兑现处理保证。
从源码实现上看,这一语义在 Function.proto 中被定义为枚举类型ProcessingGuarantees,包含三个取值:
enum ProcessingGuarantees { ATLEAST_ONCE = 0; // [default value] ATMOST_ONCE = 1; EFFECTIVELY_ONCE = 2; }其中ATLEAST_ONCE是枚举的默认值(default value),与文档中“未指定时默认ATLEAST_ONCE”的描述一一对应。该枚举同时被FunctionDetails消息引用(字段processingGuarantees,编号 6),是整个函数/连接器运行时语义的基础。Source 配置在 SourceConfigUtils.java 中通过convertProcessingGuarantee把用户配置转换为 proto 枚举,转换逻辑定义在 FunctionCommon.java。
在运行时,处理保证语义直接决定了 Source 消费消息后的确认行为。以 PulsarSource.java 为例:当处理保证为EFFECTIVELY_ONCE时,Source 使用acknowledgeCumulativeAsync(累积确认)来确认消息;否则使用acknowledgeAsync(单条确认)。而在消息处理失败时,若为EFFECTIVELY_ONCE会直接抛出运行时异常,其他语义下则对消息执行negativeAcknowledge(负向确认,触发重新投递)——这正是at-least-once语义下“消息可能被处理多次”的实现基础。
设置处理保证(Set)
创建连接器时,可以通过以下语义值设置处理保证:
ATLEAST_ONCEATMOST_ONCEEFFECTIVELY_ONCE
如果创建连接器时未指定
--processing-guarantees参数,默认语义为ATLEAST_ONCE。
以下以Admin CLI为例展示设置方法;REST API 与 Java Admin API 的对应用法可参考 io-use.md#create。
创建Source并设置处理保证:
$ bin/pulsar-admin sources create \ --processing-guarantees ATMOST_ONCE \ # 其他 source 配置项
pulsar-admin sources create的完整参数列表参见 reference-connector-admin.md#create。
创建Sink并设置处理保证:
$ bin/pulsar-admin sinks create \ --processing-guarantees EFFECTIVELY_ONCE \ # 其他 sink 配置项
pulsar-admin sinks create的完整参数列表参见 reference-connector-admin.md#create-1。
更新处理保证(Update)
连接器创建之后,可以通过以下命令更新其处理保证,取值同样为:
ATLEAST_ONCEATMOST_ONCEEFFECTIVELY_ONCE
更新Source的处理保证:
$ bin/pulsar-admin sources update \ --processing-guarantees EFFECTIVELY_ONCE \ # 其他 source 配置项
pulsar-admin sources update的完整参数列表参见 reference-connector-admin.md#update。
更新Sink的处理保证:
$ bin/pulsar-admin sinks update \ --processing-guarantees ATMOST_ONCE \ # 其他 sink 配置项
pulsar-admin sinks update的完整参数列表参见 reference-connector-admin.md#update-1。
与连接器协作:创建、运行与管理
你可以通过 Connector Admin CLI,使用 sources 和 sinks 子命令来管理 Pulsar 连接器(例如创建、更新、启动、停止、重启、重新加载、删除以及其他操作)。
使用内置连接器
Pulsar 内置了多个 内置连接器,用于在常用系统(如数据库、消息系统)之间移动数据。使用内置连接器非常简单:按照 getting-started-standalone.md 中安装内置连接器的说明 完成安装后,所有内置连接器都会被 Pulsar Broker(或 Functions Worker)自动发现,无需额外的安装步骤。
从 2.3.0 版本开始,Pulsar 将全部内置连接器作为独立的 NAR 归档发布。启用这些内置连接器,需要先从下载页面下载连接器 NAR 归档,然后将其放入解压后的 Pulsar 发行版目录下的connectors目录:
# 解压 Pulsar tarball 并将连接器 NAR 归档复制到 connectors 目录 $ tar xvfz /path/to/apache-pulsar-<version>-bin.tar.gz $ cd apache-pulsar-<version> $ mkdir connectors $ cp -r /path/to/downloaded/connectors/*.nar ./connectors $ ls connectors pulsar-io-aerospike-<version>.nar pulsar-io-cassandra-<version>.nar pulsar-io-kafka-<version>.nar pulsar-io-kinesis-<version>.nar pulsar-io-rabbitmq-<version>.nar pulsar-io-twitter-<version>.nar ...也可以直接使用自带全部内置连接器的 Docker 镜像apachepulsar/pulsar-all:<version>。
每个连接器模块的pulsar-io.yaml文件(位于META-INF/services/下)声明了连接器的类型信息。例如 pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml 中定义了:
name: cassandra description: Writes data into Cassandra sinkClass: org.apache.pulsar.io.cassandra.CassandraStringSink sinkConfigClass: org.apache.pulsar.io.cassandra.CassandraSinkConfig这也解释了文档中的一条注意事项:内置连接器的sink-type/source-type参数值由pulsar-io.yaml中的name参数决定。例如 Cassandra Sink 的--sink-type cassandra即来源于此文件中的name: cassandra。
配置连接器:YAML 配置文件
配置 Pulsar IO 连接器非常简单:运行连接器时提供一个 YAML 配置文件即可。该 YAML 配置文件告诉 Pulsar 去哪里定位 Source 和 Sink,以及如何将 Source / Sink 与 Pulsar 主题连接起来。
下面是一个 Cassandra Sink 的 YAML 配置示例:
tenant: public namespace: default name: cassandra-test-sink ... # cassandra 专用配置 configs: roots: "localhost:9042" keyspace: "pulsar_test_keyspace" columnFamily: "pulsar_test_table" keyname: "key" columnName: "col"这个示例告诉 Pulsar 要连接哪个 Cassandra 集群(roots)、在 Cassandra 中使用哪个keyspace和columnFamily来收集数据,以及如何将 Pulsar 消息映射为 Cassandra 表的 key 和列(keyname、columnName)。各连接器的详细配置请查阅 io-overview.md 中连接器清单 对应的各连接器文档。
运行 Source 连接器
可以使用如下形式的命令,将一个 Source 提交到现有 Pulsar 集群中运行:
$ ./bin/pulsar-admin sources create --classname <classname> --archive <jar-location> --tenant <tenant> --namespace <namespace> --name <source-name> --destination-topic-name <output-topic>示例:
bin/pulsar-admin sources create --classname org.apache.pulsar.io.twitter.TwitterFireHose --archive ~/application.jar --tenant test --namespace ns1 --name twitter-source --destination-topic-name twitter_data除了提交到集群运行,也可以将 Source 作为本地机器上的独立进程运行:
bin/pulsar-admin sources localrun --classname org.apache.pulsar.io.twitter.TwitterFireHose --archive ~/application.jar --tenant test --namespace ns1 --name twitter-source --destination-topic-name twitter_data如果提交的是内置 Source,则无需指定--classname和--archive,只需通过--source-type指定类型即可,命令形式如下:
./bin/pulsar-admin sources create \ --tenant <tenant> \ --namespace <namespace> \ --name <source-name> \ --destination-topic-name <input-topics> \ --source-type <source-type>以提交一个 Kafka Source 为例:
./bin/pulsar-admin sources create \ --tenant test-tenant \ --namespace test-namespace \ --name test-kafka-source \ --destination-topic-name pulsar_sink_topic \ --source-type kafka运行 Sink 连接器
可以使用如下形式的命令,将一个 Sink 提交到现有 Pulsar 集群中运行:
./bin/pulsar-admin sinks create --classname <classname> --archive <jar-location> --tenant test --namespace <namespace> --name <sink-name> --inputs <input-topics>示例:
./bin/pulsar-admin sinks create --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic同样,也可以将 Sink 作为本地进程运行:
./bin/pulsar-admin sinks localrun --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic如果提交的是内置 Sink,无需指定--classname和--archive,只需通过--sink-type指定类型:
注意:内置连接器的
sink-type参数值由pulsar-io.yaml文件中name参数的设置决定。
./bin/pulsar-admin sinks create \ --tenant <tenant> \ --namespace <namespace> \ --name <sink-name> \ --inputs <input-topics> \ --sink-type <sink-type>以提交一个 Cassandra Sink 为例:
./bin/pulsar-admin sinks create \ --tenant test-tenant \ --namespace test-namespace \ --name test-cassandra-sink \ --inputs pulsar_input_topic \ --sink-type cassandra监控连接器:查看元数据与运行状态
由于 Pulsar IO 连接器是以 Pulsar Functions 的形式运行的,因此可以使用 pulsar-admin CLI 工具中的functions命令来监控它们。
获取连接器元数据(Metadata)
bin/pulsar-admin functions get \ --tenant <tenant> \ --namespace <namespace> \ --name <connector-name>获取连接器运行状态(Status)
bin/pulsar-admin functions getstatus \ --tenant <tenant> \ --namespace <namespace> \ --name <connector-name>此外,pulsar-admin sources与pulsar-admin sinks子命令也分别提供了status、get等操作。例如在 Cassandra Sink 示例中,sink status的返回结果包含numInstances(实例总数)、numRunning(运行中实例数)以及每个实例的numReadFromPulsar(已从 Pulsar 读取的消息数)、numWrittenToSink(已写入 Sink 的消息数)等统计字段,可用于判断连接器的工作状态。
实战:连接 Pulsar 与 Apache Cassandra
下面通过一个端到端示例,演示如何使用内置 Cassandra Sink 将 Pulsar 主题中的数据写入 Cassandra,全程无需编写任何代码。完整步骤参考 io-quickstart.md。
前提:以下操作假设 Pulsar 以 standalone 模式 运行,并且所有命令都在 Pulsar 二进制发行版的根目录下执行。这些命令同样适用于多节点 Pulsar 集群。
启动 Pulsar 服务
bin/pulsar standalone所有 Pulsar 组件将按顺序启动。可以通过下面的端点检查服务是否正常:
- 检查 Pulsar 二进制协议端口(默认 6650):
telnet localhost 6650 - 检查 Pulsar Functions 集群:
curl -s http://localhost:8080/admin/v2/worker/cluster - 确认
publictenant 与defaultnamespace 存在:curl -s http://localhost:8080/admin/v2/namespaces/public - 确认所有内置连接器均可用:
curl -s http://localhost:8080/admin/v2/functions/connectors
其中第 4 步的示例输出(JSON 数组)列出了每个内置连接器的名称、描述以及对应的 Source/Sink 实现类,例如:
[{"name":"aerospike","description":"Aerospike database sink","sinkClass":"org.apache.pulsar.io.aerospike.AerospikeStringSink"},{"name":"cassandra","description":"Writes data into Cassandra","sinkClass":"org.apache.pulsar.io.cassandra.CassandraStringSink"},{"name":"kafka","description":"Kafka source and sink connector","sourceClass":"org.apache.pulsar.io.kafka.KafkaStringSource","sinkClass":"org.apache.pulsar.io.kafka.KafkaBytesSink"},{"name":"kinesis","description":"Kinesis sink connector","sinkClass":"org.apache.pulsar.io.kinesis.KinesisSink"},{"name":"rabbitmq","description":"RabbitMQ source connector","sourceClass":"org.apache.pulsar.io.rabbitmq.RabbitMQSource"},{"name":"twitter","description":"Ingest data from Twitter firehose","sourceClass":"org.apache.pulsar.io.twitter.TwitterFireHose"}]如果启动过程中出错,可以在运行pulsar standalone的终端看到异常,也可以查看 Pulsar 目录下logs目录中的日志。
启动单节点 Cassandra 集群
使用 Docker 启动一个单节点 Cassandra 集群:
docker run -d --rm --name=cassandra -p 9042:9042 cassandra启动后依次验证:docker ps确认进程在运行;docker logs cassandra检查日志;docker exec cassandra nodetool status检查集群状态(输出中节点应显示为UN,即 Up/Normal)。
使用cqlsh连接并创建 keyspace 与表:
$ docker exec -ti cassandra cqlsh localhost Connected to Test Cluster at localhost:9042. [cqlsh 5.0.1 | Cassandra 3.11.2 | CQL spec 3.4.4 | Native protocol v4]cqlsh> CREATE KEYSPACE pulsar_test_keyspace WITH replication = {'class':'SimpleStrategy', 'replication_factor':1}; cqlsh> USE pulsar_test_keyspace; cqlsh:pulsar_test_keyspace> CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);配置并提交 Cassandra Sink
创建配置文件examples/cassandra-sink.yml:
configs: roots: "localhost:9042" keyspace: "pulsar_test_keyspace" columnFamily: "pulsar_test_table" keyname: "key" columnName: "col"提交 Cassandra Sink(sink-type取自 pulsar-io.yaml 中的name):
bin/pulsar-admin sink create \ --tenant public \ --namespace default \ --name cassandra-test-sink \ --sink-type cassandra \ --sink-config-file examples/cassandra-sink.yml \ --inputs test_cassandra命令执行后,Pulsar 会创建一个名为cassandra-test-sink的 Sink 连接器,它将以 Pulsar Function 的形式运行,并把主题test_cassandra中产生的消息写入 Cassandra 表pulsar_test_table。
查看 Sink 信息与状态
bin/pulsar-admin sink get \ --tenant public \ --namespace default \ --name cassandra-test-sink示例输出:
{ "tenant": "public", "namespace": "default", "name": "cassandra-test-sink", "className": "org.apache.pulsar.io.cassandra.CassandraStringSink", "inputSpecs": { "test_cassandra": { "isRegexPattern": false } }, "configs": { "roots": "localhost:9042", "keyspace": "pulsar_test_keyspace", "columnFamily": "pulsar_test_table", "keyname": "key", "columnName": "col" }, "parallelism": 1, "processingGuarantees": "ATLEAST_ONCE", "retainOrdering": false, "autoAck": true, "archive": "builtin://cassandra" }可以看到,未显式指定处理保证时,processingGuarantees默认即为ATLEAST_ONCE,与本文前面介绍的默认语义完全一致;archive: builtin://cassandra说明这是一个内置连接器。
检查运行状态:
bin/pulsar-admin sink status \ --tenant public \ --namespace default \ --name cassandra-test-sink生产消息并验证数据落库
向 Sink 的输入主题test_cassandra生产 10 条消息:
for i in {0..9}; do bin/pulsar-client produce -m "key-$i" -n 1 test_cassandra; done再次查看 Sink 状态,可以看到numReadFromPulsar与numWrittenToSink均为 10:
{ "numInstances" : 1, "numRunning" : 1, "instances" : [ { "instanceId" : 0, "status" : { "running" : true, "error" : "", "numRestarts" : 0, "numReadFromPulsar" : 10, "numSystemExceptions" : 0, "latestSystemExceptions" : [ ], "numSinkExceptions" : 0, "latestSinkExceptions" : [ ], "numWrittenToSink" : 10, "lastReceivedTime" : 1551685489136, "workerId" : "c-standalone-fw-localhost-8080" } } ] }最后在 Cassandra 中验证数据:
docker exec -ti cassandra cqlsh localhostcqlsh> use pulsar_test_keyspace; cqlsh:pulsar_test_keyspace> select * from pulsar_test_table; key | col --------+-------- key-5 | key-5 key-0 | key-0 ...删除 Sink
bin/pulsar-admin sink delete \ --tenant public \ --namespace default \ --name cassandra-test-sink从源码理解连接器的实现基础
- 运行时模型:连接器本质上以 Function 的形态运行。定义连接器运行细节的核心模型位于 Function.proto,其中
SourceSpec与SinkSpec消息分别承载 Source / Sink 的类名、配置(JSON 格式的configs)、内置连接器标识builtin、订阅类型等字段,FunctionDetails统一管理租户、命名空间、名称、并行度(parallelism)与处理保证(processingGuarantees)。 - 处理保证的落点:Source 的配置转换在 SourceConfigUtils.java 中完成,Sink 对应 SinkConfigUtils.java;两者通过 FunctionCommon.java 中的
convertProcessingGuarantee在用户配置与 proto 枚举之间互转。运行时,PulsarSource.java 依据处理保证选择累积确认或单条确认,PulsarSink.java 负责将消息写入外部系统。 - 测试验证:仓库中 SourceConfigUtilsTest.java 与 SinkConfigUtilsTest.java 覆盖了配置转换逻辑,PulsarSourceTest.java 与 PulsarSinkTest.java 覆盖了运行时行为,可作为深入阅读与二次开发的入口。
总结
Pulsar IO 连接器通过 Source / Sink 两种抽象,打通了 Pulsar 与外部数据系统之间的双向数据通道;其处理保证(at-most-once、at-least-once、effectively-once)与 Pulsar Functions 保持统一,默认语义为ATLEAST_ONCE,可通过pulsar-admin sources/sinks create与update命令设置或更新。内置连接器以 NAR 归档形式发布,被 Functions Worker 自动发现;连接器以 Function 实例的形态运行在 Worker 上,既可用sources/sinks子命令管理,也可用functions子命令监控。结合本仓库源码,读者可以进一步深入理解处理保证在消费确认、失败重试等环节的具体落点,从而在生产环境中正确选择与配置连接器。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar IO Connectors 全景指南:Source 与 Sink 架构、处理语义与 Connector Admin 实操
Apache Pulsar IO Connectors 全景指南:Source 与 Sink 架构、处理语义与 Connector Admin 实操 Apach
消息队列后端流处理Apache Pulsar IO 连接器详解:Source、Sink 架构与处理语义配置实战
Apache Pulsar IO 连接器详解:Source、Sink 架构与处理语义配置实战 Pulsar IO 是 Apache Pulsar 内置的流数据接
消息队列后端流处理Apache Pulsar IO 连接器概览:Source/Sink 架构与处理语义(Processing Guarantees)深度解析
Apache Pulsar IO 连接器概览:Source/Sink 架构与处理语义(Processing Guarantees)深度解析 本文是 Apache
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考