- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文以 Apache Pulsar IO(Pulsar IO framework)为背景,系统讲解如何在 Pulsar 集群中部署、配置、运行、监控与升级 Connector(Source / Sink)。你将掌握pulsar-adminCLI 的source/sink命令族、YAML 配置文件的编写方式、内置连接器的自动发现机制,以及如何通过 Pulsar Functions 命令获取连接器元数据与运行状态,从而在生产环境中完整地管理数据进出 Pulsar 的管道。
本文对应文档版本为 Pulsar 2.2.0,位于 site2/website-next/versioned_docs/version-2.2.0/io-managing.md;仓库中较新的 CLI 文档 site2/website-next/docs/io-cli.md 对命令族有更完整的参数说明,本文在继承原文档骨架的同时结合两者与源码进行深度扩充。
一、Pulsar IO 连接器管理概览
Pulsar IO 是 Apache Pulsar 内置的数据集成框架,它把外部系统与 Pulsar 之间的数据流动抽象为两类连接器:
- Source(数据源):从外部系统(如 Kafka、Twitter、文件、Netty 等)读取数据,写入 Pulsar 主题,实现数据"入"Pulsar(ingress)。
- Sink(数据汇):从 Pulsar 主题消费消息,写入外部系统(如 Cassandra、HBase、Elasticsearch、JDBC 等),实现数据"出"Pulsar(egress)。
连接器的核心价值在于:你不需要自己写消费者或生产者代码,只需通过一个 YAML 配置文件描述"连接到哪里、如何映射",再通过pulsar-admin命令行把连接器提交到集群即可。原文档将管理任务归纳为四类,本文逐一展开:
- 部署内置连接器(Deploy builtin connectors)
- 用 Pulsar Admin CLI 监控和更新运行中的连接器(Monitor and update running connectors)
- 部署自定义连接器(Deploy customized connectors)
- 升级连接器(Upgrade a connector)
二、使用内置连接器
Pulsar 随发行版捆绑了一批内置连接器(builtin connectors),用于与常见的数据库、消息系统等双向搬运数据。完整清单可参阅 site2/website-next/docs/io-overview.md 中的 "Working with connectors" 一节(该文档目录下还有每个连接器的独立说明,例如 Cassandra 见 site2/website-next/docs/io-cassandra.md、Kafka 见 site2/website-next/docs/io-kafka.md)。
2.1 安装内置连接器
内置连接器的安装步骤见 site2/website-next/docs/getting-started-standalone.md 中 "Installing builtin connectors" 一节。安装完成后,所有内置连接器会被 Pulsar Broker(或 Function Worker)自动发现,无需额外的安装步骤。
2.2 自动发现机制:pulsar-io.yaml
"自动发现"并非魔法,而是源于每个连接器 NAR 包内携带的META-INF/services/pulsar-io.yaml描述文件。仓库中每个连接器模块都包含这样一个文件,例如 Cassandra 连接器:
- 描述文件: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这个文件声明了三件关键信息:连接器对外暴露的name(即 CLI 中的--sink-type/--source-type)、连接器实现类sinkClass、以及配置类sinkConfigClass。原文档特别强调:内置连接器的sink-type参数由pulsar-io.yaml文件中的name字段决定——也就是说,你在 CLI 里填写的cassandra正是这个name值。仓库中pulsar-io/目录下的 30 余个模块(aerospike、kafka、kinesis、rabbitmq、redis、hdfs2、hdfs3、jdbc/*等)各自带有同名描述文件,共同构成内置连接器注册表。
从源码结构看,Broker/Function Worker 在启动时会扫描各 NAR 包中的pulsar-io.yaml,将name与实现类建立映射,因此pulsar-admin sources available-sources、pulsar-admin sinks available-sinks可以直接枚举出集群当前支持的所有内置连接器(见下文监控小节)。
三、配置连接器:以 Cassandra Sink 为例
配置 Pulsar IO 连接器非常直接:在运行连接器(Running Connectors)时提供一个 YAML 配置文件。该 YAML 告诉 Pulsar 三件事:Source/Sink 实现类位于哪里(archive/classname)、连接器归属哪个租户与命名空间、以及如何把外部系统与 Pulsar 主题对接。
原文档给出的 Cassandra Sink 配置示例:
tenant: public namespace: default name: cassandra-test-sink ... # cassandra specific config configs: roots: "localhost:9042" keyspace: "pulsar_test_keyspace" columnFamily: "pulsar_test_table" keyname: "key" columnName: "col"这个示例的含义是:Pulsar 连接到哪个 Cassandra 集群(roots)、数据落到哪个keyspace与columnFamily、以及如何把一条 Pulsar 消息映射为 Cassandra 表的主键(keyname)与列(columnName)。
3.1 源码级字段说明
以上configs段落的字段并非随意约定,它们与配置类 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java 一一对应,且全部标注为required = true:
| YAML 字段 | 源码字段 | 类型 | 说明 |
|---|---|---|---|
roots | roots | String | 一组以逗号分隔的 Cassandra 主机列表(如localhost:9042),必填 |
keyspace | keyspace | String | 用于写入 Pulsar 消息的 keyspace,必填 |
columnFamily | columnFamily | String | Cassandra 列族(表)名称,必填 |
keyname | keyname | String | Cassandra 列族中作为主键的列名,必填 |
columnName | columnName | String | Cassandra 列族中用于写入消息内容的列名,必填 |
该配置类的load(String yamlFile)方法使用 Jackson 的 YAML 工厂(ObjectMapper(new YAMLFactory()))将上述 YAML 反序列化为配置对象,这也是"提供 YAML 即完成配置"这一机制的实现入口。
3.2 roots 的解析细节
在 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java 的createClient(String roots)方法中,roots会按逗号切分成多个主机,逐个加入Cluster.builder():
String[] hosts = roots.split(","); if (hosts.length <= 0) { throw new RuntimeException("Invalid cassandra roots"); }这意味着你可以在roots中配置多个 Cassandra 节点以实现连接层面的容错;若配置为空或格式非法,Sink 会在启动阶段直接抛出Invalid cassandra roots异常。
提示:
configs段落的字段随连接器类型而异,具体字段请查阅 site2/website-next/docs/io-overview.md 中每个独立连接器的文档页。
四、运行连接器
连接器通过pulsar-adminCLI 工具的source与sink命令族管理,pulsar-admin的完整命令参考见 site2/website-next/docs/reference-pulsar-admin.md。
版本说明:2.2.0 文档中的命令名为
source create/sink create(单数形式);仓库新版本 CLI 文档 site2/website-next/docs/io-cli.md 已统一为sources create/sinks create(复数形式),功能等价,下文参数说明以新版文档为准。
4.1 运行 Source(数据源)
方式一:提交到集群运行(cluster 模式)
使用以下形式把自定义 Source 提交到已有 Pulsar 集群:
$ ./bin/pulsar-admin source create --classname <classname> --archive <jar-location> --tenant <tenant> --namespace <namespace> --name <source-name> --destination-topic-name <output-topic>示例(提交 Twitter Firehose Source):
bin/pulsar-admin source create --classname org.apache.pulsar.io.twitter.TwitterFireHose --archive ~/application.jar --tenant test --namespace ns1 --name twitter-source --destination-topic-name twitter_data方式二:本地进程运行(localrun 模式)
如果不希望把 Source 提交到集群,也可以让它在本地机器上以独立进程运行:
bin/pulsar-admin source localrun --classname org.apache.pulsar.io.twitter.TwitterFireHose --archive ~/application.jar --tenant test --namespace ns1 --name twitter-source --destination-topic-name twitter_datalocalrun模式非常适合开发调试阶段:它不经过 Function Worker 调度,而是把连接器直接跑在当前机器的 JVM 进程中,并可通过--broker-service-url指定要连接的 Broker 地址。
方式三:提交内置 Source(免 classname/archive)
如果提交的是内置 Source,则无需指定--classname与--archive,只需给出--source-type:
./bin/pulsar-admin source create \ --tenant <tenant> \ --namespace <namespace> \ --name <source-name> \ --destination-topic-name <input-topics> \ --source-type <source-type>示例(提交 Kafka Source):
./bin/pulsar-admin source create \ --tenant test-tenant \ --namespace test-namespace \ --name test-kafka-source \ --destination-topic-name pulsar_sink_topic \ --source-type kafka--source-type kafka对应 Kafka 连接器 NAR 中pulsar-io.yaml的name字段(参见 pulsar-io/kafka/src/main/resources/META-INF/services/pulsar-io.yaml),Broker 据此定位实现类并完成加载。
4.2 运行 Sink(数据汇)
方式一:提交到集群运行(cluster 模式)
./bin/pulsar-admin sink create --classname <classname> --archive <jar-location> --tenant test --namespace <namespace> --name <sink-name> --inputs <input-topics>示例(提交 Cassandra Sink):
./bin/pulsar-admin sink create --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic注意示例中的--classname org.apache.pulsar.io.cassandra是文档原样给出的简写形式;实际以 NAR 描述文件为准,Cassandra Sink 的实现类完整名称为org.apache.pulsar.io.cassandra.CassandraStringSink。
方式二:本地进程运行(localrun 模式)
./bin/pulsar-admin sink localrun --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic方式三:提交内置 Sink(免 classname/archive)
提交内置 Sink 时同样只需--sink-type:
./bin/pulsar-admin sink create \ --tenant <tenant> \ --namespace <namespace> \ --name <sink-name> \ --inputs <input-topics> \ --sink-type <sink-type>注意:内置连接器的
sink-type参数由pulsar-io.yaml文件中的name参数决定(如前文 Cassandra 描述文件中的name: cassandra)。
示例(提交 Cassandra Sink):
./bin/pulsar-admin sink create \ --tenant test-tenant \ --namespace test-namespace \ --name test-cassandra-sink \ --inputs pulsar_input_topic \ --sink-type cassandra4.3 关键参数速查(来自仓库 CLI 文档)
根据 site2/website-next/docs/io-cli.md,sources/sinks两个命令族均包含 12 个子命令:create、update、delete、get、status、list、stop、start、restart、localrun、available-sources(或available-sinks)、reload。核心参数整理如下:
| 参数 | 适用命令 | 说明 |
|---|---|---|
-a, --archive | create/update/localrun | NAR 归档包路径,也支持 http/https/file URL(file 协议要求包已存在于 Worker 主机上) |
--classname | create/update/localrun | archive 为 file:// 路径时的连接器实现类名 |
-t, --source-type / --sink-type | create/update | 内置连接器类型,由pulsar-io.yaml的name决定 |
--destination-topic-name | source | Source 输出数据写入的 Pulsar 主题 |
-i, --inputs | sink | Sink 消费的输入主题(多个可用逗号分隔) |
--tenant/--namespace/--name | 全部 | 连接器的租户、命名空间与名称,唯一标识一个连接器 |
--parallelism | create/update | 并行度,即运行多少个连接器实例 |
--processing-guarantees | create/update | 处理语义:ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE |
--source-config-file / --sink-config-file | create/update | 上文第三节所述 YAML 配置文件的路径 |
--source-config / --sink-config | create/update | 以 key/value 形式直接传入配置 |
--cpu/--ram/--disk | create/update | 每个实例分配的资源(CPU 核数、RAM 字节数、磁盘字节数) |
--retain-ordering | sink | 是否按序消费并写入消息 |
--timeout-ms | sink | 消息超时时间(毫秒) |
--subs-name | sink | 输入主题消费时使用的订阅名称 |
--broker-service-url | localrun | 本地运行时连接的 Broker 地址 |
--tls-*、--client-auth-* | localrun | TLS 与客户端认证相关参数 |
这些选项在源码中的解析入口位于 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java 与 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java,感兴趣的读者可进一步追踪参数到 Pulsar Functions 运行时配置对象的映射过程。
五、监控连接器
由于Pulsar IO 连接器本质上以 Pulsar Functions 的形式运行(这一点原文档明确指出,Pulsar Functions 概览见 site2/website-next/docs/functions-overview.md),因此你可以直接使用pulsar-admin的functions命令来监控它们。
5.1 获取连接器元数据
bin/pulsar-admin functions get \ --tenant <tenant> \ --namespace <namespace> \ --name <connector-name>functions get返回连接器的完整元数据,包括实现类、并行度、配置、输入/输出主题、处理语义等,等同于查看提交时登记的状态快照,是排查"连接器配置是否正确"的首选命令。
5.2 获取连接器运行状态
bin/pulsar-admin functions getstatus \ --tenant <tenant> \ --namespace <namespace> \ --name <connector-name>functions getstatus返回连接器各实例的实时运行状态(如每个实例的接收/处理消息数、最近错误等),用于判断连接器是否健康、是否有背压或异常堆积。
补充:除
functions命令外,sources status与sinks status也可查看对应连接器的运行状态(--instance-id可指定实例,缺省时返回全部实例),sources list/sinks list可列出某租户命名空间下所有运行中的连接器;sources available-sources/sinks available-sinks可枚举集群支持的内置连接器,sources reload/sinks reload可重新加载内置连接器注册表(适用于新增 NAR 后无需重启即可感知的场景)。完整用法见 site2/website-next/docs/io-cli.md。
六、更新、升级与删除连接器
原文档开篇提出的目标中包含"更新运行中的连接器"与"升级连接器",这两项分别对应sources update/sinks update与delete子命令:
6.1 更新连接器配置或实现
update子命令用于修改已提交连接器的参数,用法与create相同,可更新的内容包括--archive、--classname、--parallelism、--processing-guarantees、--source-config-file/--sink-config-file等:
# 以更新 Sink 为例 bin/pulsar-admin sink update \ --tenant test-tenant \ --namespace test-namespace \ --name test-cassandra-sink \ --sink-config-file /path/to/new-config.yaml \ --parallelism 4update是滚动升级连接器实现(例如更换 NAR 版本或调整并行度)的主要途径;Sink 的 update 还额外支持--update-auth-data参数(默认false)决定是否同时更新认证数据。
6.2 删除连接器
bin/pulsar-admin sink delete \ --tenant <tenant> \ --namespace <namespace> \ --name <sink-name>Source 的删除命令与之同理(source delete)。删除后,由该连接器创建的相关运行实例与调度状态会被清理。
6.3 停止 / 启动 / 重启实例
针对排查问题场景,连接器实例可被单独或整体控制:
sources stop/sinks stop:停止指定实例(--instance-id),缺省停止全部实例;sources start/sinks start:启动实例;sources restart/sinks restart:重启实例。
这三个子命令在"连接器卡死需要临时摘流"或"修改外部系统后需要重连"等运维场景中非常实用。
七、完整管理流程小结
以部署一个 Cassandra Sink 为例,完整生命周期为:
- 准备配置:编写 YAML(tenant/namespace/name/configs),字段对照 CassandraSinkConfig.java;
- 提交运行:
bin/pulsar-admin sink create --sink-type cassandra --tenant ... --namespace ... --name ... --inputs ... --sink-config-file config.yaml(或本地调试用sink localrun); - 监控验证:
bin/pulsar-admin functions getstatus --tenant ... --namespace ... --name cassandra-sink观察消息处理情况; - 升级调整:配置变化用
sink update,异常时用stop/start/restart,彻底下线用sink delete。
本文覆盖了 site2/website-next/versioned_docs/version-2.2.0/io-managing.md 的全部内容(内置连接器部署、YAML 配置、Source/Sink 运行、Functions 监控),并以仓库中 site2/website-next/docs/io-cli.md 与pulsar-io/cassandra模块源码为佐证,补充了参数表、pulsar-io.yaml自动发现机制、配置类字段说明与更新/删除/重启等运维细节。实践中请以你所使用 Pulsar 版本的官方文档为准,并注意本仓库对应 Pulsar 2.2.0 文档中source/sink单数命令与新版sources/sinks复数命令的差异。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 连接器管理实战:内置连接器部署、配置、运行与监控
Apache Pulsar 连接器管理实战:内置连接器部署、配置、运行与监控 本篇技术指南聚焦 Apache Pulsar 的 IO Connectors(连接
消息队列后端流处理Apache Pulsar IO 连接器管理实战:内置连接器的部署、运行、监控与升级(Pulsar 2.3.0)
Apache Pulsar IO 连接器管理实战:内置连接器的部署、运行、监控与升级(Pulsar 2.3.0) 本篇指南聚焦 Apache Pulsar IO
消息队列后端流处理Apache Pulsar 内置连接器(Pulsar IO Connectors):安装、配置与实战详解
Apache Pulsar 内置连接器(Pulsar IO Connectors):安装、配置与实战详解 Apache Pulsar 发行版内置了一组经过打包与
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考