消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文以 Apache Pulsar IOPulsar IO framework为背景系统讲解如何在 Pulsar 集群中部署、配置、运行、监控与升级 ConnectorSource / 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 主题实现数据入Pulsaringress。Sink数据汇从 Pulsar 主题消费消息写入外部系统如 Cassandra、HBase、Elasticsearch、JDBC 等实现数据出Pulsaregress。连接器的核心价值在于你不需要自己写消费者或生产者代码只需通过一个 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.yamlname: 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 trueYAML 字段源码字段类型说明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 Sourcebin/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.yamlBroker 据此定位实现类并完成加载。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.mdsources/sinks两个命令族均包含 12 个子命令create、update、delete、get、status、list、stop、start、restart、localrun、available-sources或available-sinks、reload。核心参数整理如下参数适用命令说明-a, --archivecreate/update/localrunNAR 归档包路径也支持 http/https/file URLfile 协议要求包已存在于 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-namefunctions get返回连接器的完整元数据包括实现类、并行度、配置、输入/输出主题、处理语义等等同于查看提交时登记的状态快照是排查连接器配置是否正确的首选命令。5.2 获取连接器运行状态bin/pulsar-admin functions getstatus \ --tenant tenant \ --namespace namespace \ --name connector-namefunctions 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-nameSource 的删除命令与之同理source delete。删除后由该连接器创建的相关运行实例与调度状态会被清理。6.3 停止 / 启动 / 重启实例针对排查问题场景连接器实例可被单独或整体控制sources stop/sinks stop停止指定实例--instance-id缺省停止全部实例sources start/sinks start启动实例sources restart/sinks restart重启实例。这三个子命令在连接器卡死需要临时摘流或修改外部系统后需要重连等运维场景中非常实用。七、完整管理流程小结以部署一个 Cassandra Sink 为例完整生命周期为准备配置编写 YAMLtenant/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复数命令的差异。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 连接器管理实战内置连接器部署、配置、运行与监控Apache Pulsar 连接器管理实战内置连接器部署、配置、运行与监控 本篇技术指南聚焦 Apache Pulsar 的 IO Connectors连接消息队列后端流处理Apache Pulsar IO 连接器管理实战内置连接器的部署、运行、监控与升级Pulsar 2.3.0Apache Pulsar IO 连接器管理实战内置连接器的部署、运行、监控与升级Pulsar 2.3.0 本篇指南聚焦 Apache Pulsar IO消息队列后端流处理Apache Pulsar 内置连接器Pulsar IO Connectors安装、配置与实战详解Apache Pulsar 内置连接器Pulsar IO Connectors安装、配置与实战详解 Apache Pulsar 发行版内置了一组经过打包与消息队列后端流处理上一篇prek 实测Rust 重写的 Git 钩子管理器比 pre-commit 快近 13 倍下一篇SuckIT实战案例如何快速下载静态网站、博客和文档站点到本地创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
Learn-Algorithms 链表删除专题:O(1) 时间删除结点与双向循环链表去重的实战解析 教程 【免费下载链接】Learn-Algorithms 算法学习笔记 项目地址: https://gitcode.com/gh_mirrors/le/Learn-Algorithms 点击查看 免费下载 链表删除是算法面试中的高频考点。本专题基于仓库笔记 2.2 链表-删除.md,系统梳理四类删除问题:在 … · 2026/9/25 7:09:36
PADS Logic原理图设计:工程配置、网络标签与DRC验证实践 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:09:36
认识 React 360:用 React 构建跨平台 360° 与 VR 网页应用 前端3D渲染 【免费下载链接】react-360 Create amazing 360 and VR content using React 项目地址: https://gitcode.com/gh_mirrors/re/react-360 点击查看 免费下载 React 360 是一个基于 React 构建 3D 与 VR 用户界面的开源框架,让你用熟悉的组件、… · 2026/9/25 7:09:29
19个免费PPT网站实测:在线编辑、模板下载与AI辅助工具推荐 1. 为什么我花了两周时间实测这19个PPT网站做PPT这件事,说大不大,说小也绝对不小。我在一家中型企业做品牌策划,平均每个月要出4到6份对外提案,加上内部汇报、季度复盘、培训课件,一年下来经手的PPT少说也有七八十份。… · 2026/9/25 7:33:26
Word尾注脚注管理全攻略:插入、删除与去横线技巧 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:33:26
Simulink建模效率:自动整理连线、显示数据类型与内容自适应 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:33:26
零成本监控回放方案:旧摄像头+树莓派+夸克网盘 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:33:20
彻底关闭OfficePlus:从加载项禁用、注册表修改到完全卸载的完整指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:33:20
SUMO交通仿真入门:从零搭建微观交通场景的核心指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 7:33:20
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37