☰
Apache Pulsar IO 连接器管理实战:内置连接器部署、配置、运行与监控
2026/9/25 7:09:30 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

本文以 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 字段源码字段类型说明
rootsrootsString一组以逗号分隔的 Cassandra 主机列表(如localhost:9042),必填
keyspacekeyspaceString用于写入 Pulsar 消息的 keyspace,必填
columnFamilycolumnFamilyStringCassandra 列族(表)名称,必填
keynamekeynameStringCassandra 列族中作为主键的列名,必填
columnNamecolumnNameStringCassandra 列族中用于写入消息内容的列名,必填

该配置类的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_data

localrun模式非常适合开发调试阶段:它不经过 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 cassandra

4.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, --archivecreate/update/localrunNAR 归档包路径,也支持 http/https/file URL(file 协议要求包已存在于 Worker 主机上)
--classnamecreate/update/localrunarchive 为 file:// 路径时的连接器实现类名
-t, --source-type / --sink-typecreate/update内置连接器类型,由pulsar-io.yaml的name决定
--destination-topic-namesourceSource 输出数据写入的 Pulsar 主题
-i, --inputssinkSink 消费的输入主题(多个可用逗号分隔)
--tenant/--namespace/--name全部连接器的租户、命名空间与名称,唯一标识一个连接器
--parallelismcreate/update并行度,即运行多少个连接器实例
--processing-guaranteescreate/update处理语义:ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE
--source-config-file / --sink-config-filecreate/update上文第三节所述 YAML 配置文件的路径
--source-config / --sink-configcreate/update以 key/value 形式直接传入配置
--cpu/--ram/--diskcreate/update每个实例分配的资源(CPU 核数、RAM 字节数、磁盘字节数)
--retain-orderingsink是否按序消费并写入消息
--timeout-mssink消息超时时间(毫秒)
--subs-namesink输入主题消费时使用的订阅名称
--broker-service-urllocalrun本地运行时连接的 Broker 地址
--tls-*、--client-auth-*localrunTLS 与客户端认证相关参数

这些选项在源码中的解析入口位于 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 4

update是滚动升级连接器实现(例如更换 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 为例,完整生命周期为:

  1. 准备配置:编写 YAML(tenant/namespace/name/configs),字段对照 CassandraSinkConfig.java;
  2. 提交运行:bin/pulsar-admin sink create --sink-type cassandra --tenant ... --namespace ... --name ... --inputs ... --sink-config-file config.yaml(或本地调试用sink localrun);
  3. 监控验证:bin/pulsar-admin functions getstatus --tenant ... --namespace ... --name cassandra-sink观察消息处理情况;
  4. 升级调整:配置变化用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

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

相关推荐

上一篇:prek 实测:Rust 重写的 Git 钩子管理器,比 pre-commit 快近 13 倍
下一篇:SuckIT实战案例:如何快速下载静态网站、博客和文档站点到本地

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询