- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本篇技术指南以 Pulsar 2.0 这一重大版本为核心,系统讲解其两大核心变更:将property(属性)术语全面替换为tenant(租户),以及重构 Topic 名称体系(移除 cluster 组件、引入默认值驱动的灵活命名)。读完本文,你将掌握 Pulsar 2.0 及之后版本中 Topic 完整名称与简写名称的解析规则、pulsar-admin tenants命令行管理方法,以及 2.0 引入的 Pulsar Functions 轻量计算能力的核心编程模型。
Pulsar 2.0:一次带来大胆变更的重大版本
Pulsar 2.0 是 Pulsar 平台的一个重大版本发布,为整个系统引入了一些具有突破性的变更,主要包括:
| 新特性 | 说明 |
|---|---|
| Pulsar Functions | 为 Pulsar 提供的轻量级计算(compute)能力 |
同时,2.0 版本带来了若干重大变更(Major changes),这些变更会显著影响用户的日常使用习惯,需要特别关注:
property术语被tenant取代:pulsar-admin properties命令行接口被替换为pulsar-admin tenants;- Topic 名称体系重构:名称中不再包含 cluster 组件,
property更名为tenant,并引入了基于默认值的灵活简写命名规则; /不允许出现在 Topic 名称中:本地名称(local name)部分不能包含斜杠。
下文将逐一深入这些变更的细节、动机与底层实现。
Pulsar Functions:2.0 引入的轻量计算能力
作为 Pulsar 2.0 最重要的新增特性,Pulsar Functions 是一种轻量级计算进程,它能够:
- 从一个或多个 Pulsar Topic 消费消息;
- 对每条消息应用用户提供的处理逻辑;
- 将计算结果发布到另一个 Topic。
编程模型
Pulsar Functions 的核心编程模型非常简单:函数从输入 Topic接收消息,每收到一条消息,函数会完成以下任务之一或多个:
- 对输入应用处理逻辑,并将输出写入 Pulsar 的输出 Topic或 Apache BookKeeper 状态存储;
- 将日志写入日志 Topic(便于调试);
- 递增一个计数器。
例如,可以构建如下处理链路:Python 函数监听raw-sentencesTopic 并清洗字符串(去除多余空白、转小写),随后将结果发布到sanitized-sentencesTopic;Java 函数监听sanitized-sentencesTopic 统计窗口内单词出现次数并发布到resultsTopic;最终由 Python 函数将结果写入 MySQL 表。
Word Count 示例:Java 实现
使用 Java 版 Pulsar Functions SDK 实现经典的词频统计:
package org.example.functions; import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.api.Function; import java.util.Arrays; public class WordCountFunction implements Function<String, Void> { // 每当输入 Topic 有新消息发布时,该函数被调用 @Override public Void process(String input, Context context) throws Exception { Arrays.asList(input.split(" ")).forEach(word -> { String counterKey = word.toLowerCase(); context.incrCounter(counterKey, 1); }); return null; } }上述代码中,context.incrCounter(key, amount)是 functions-develop.md 中定义的上下文计数器接口,函数可以通过它维护跨消息的累计状态。打包成 JAR 后,即可用命令行部署到集群:
$ bin/pulsar-admin functions create \ --jar target/my-jar-with-dependencies.jar \ --classname org.example.functions.WordCountFunction \ --tenant public \ --namespace default \ --name word-count \ --inputs persistent://public/default/sentences \ --output persistent://public/default/count注意这里的--tenant public --namespace default与下文将介绍的 Pulsar 2.0 默认租户/命名空间体系完全对应。
消息处理语义
Pulsar Functions 提供三种消息投递语义,可在创建函数时通过--processing-guarantees指定:
| 投递语义 | 说明 |
|---|---|
| At-most-once(至多一次) | 每条消息可能被处理,也可能不被处理("至多"一次) |
| At-least-once(至少一次) | 每条消息可能被处理多次("至少"一次) |
| Effectively-once(恰好一次/有效一次) | 每条消息对应一个输出结果 |
例如,下面的命令创建一个EFFECTIVELY_ONCE语义的函数:
$ bin/pulsar-admin functions create \ --name my-effectively-once-function \ --processing-guarantees EFFECTIVELY_ONCE # 其他函数配置可选值为ATMOST_ONCE、ATLEAST_ONCE、EFFECTIVELY_ONCE。默认情况下,不指定该参数时函数提供 at-least-once 投递保证。创建后还可通过pulsar-admin functions update --processing-guarantees ATMOST_ONCE动态调整语义。
重大变更一:Property 更名为 Tenant
在 Pulsar 2.0 之前,系统使用property(属性)这一概念进行多租户隔离。而 property 本质上与 tenant(租户)是同一个东西,因此在 2.0 版本中,"property" 术语被正式移除,统一改称为tenant。
这一变更最直观的体现就是命令行接口的替换:
- 原
pulsar-admin properties接口 → 现pulsar-admin tenants接口。
在 reference-pulsar-admin.md 中可以看到,tenants命令组提供完整的租户生命周期管理操作:
$ pulsar-admin tenants subcommand其子命令包括:
list:列出已有租户 ——pulsar-admin tenants listget:获取指定租户配置 ——pulsar-admin tenants get tenant-namecreate:创建新租户 ——pulsar-admin tenants create tenant-name optionsupdate:更新租户 ——pulsar-admin tenants update tenant-name optionsdelete:删除租户 ——pulsar-admin tenants delete tenant-name
其中create与update支持以下选项:
| 选项 | 说明 | 默认值 |
|---|---|---|
-r,--admin-roles | 逗号分隔的管理员角色列表 | 无 |
-c,--allowed-clusters | 逗号分隔的允许访问的集群列表 | 无 |
而delete支持-f, --force选项,用于强制删除租户(同时删除其下所有 namespace),默认值为false。
从源码结构看,这一 CLI 实现位于 CmdTenants.java,其内部通过getAdmin().tenants()调用管理 API 完成getTenants()、getTenantInfo()、createTenant()、updateTenant()、deleteTenant()等操作,对应了文档所述的子命令集合。
注意:在某些场景下 "properties" 术语仍然被使用,但现在已视为**废弃(deprecated)**用法,将在未来版本中被彻底移除。
重大变更二:Topic 名称体系重塑
旧命名格式
在 Pulsar 2.0 之前,所有Pulsar Topic 的名称都采用如下形式:
{persistent|non-persistent}://property/cluster/namespace/topic即名称由四个部分组成:持久化类型、property(属性)、cluster(集群)、namespace(命名空间)以及 topic 名称本身。
新命名格式与三项关键变化
Pulsar 2.0 对 Topic 名称做出了以下重要调整:
- 移除了 cluster 组件(详见下文);
- property 更名为 tenant;
- 引入了灵活的简写命名体系,可大幅缩短大部分 Topic 名称;
/不允许出现在 Topic 名称中(作为分隔符之外的本地名称部分不允许含斜杠)。
移除 cluster 组件后,所有 Topic 名称现在都采用如下形式:
{persistent|non-persistent}://tenant/namespace/topic即从四段式简化为三段式:tenant/namespace/topic。
在 TopicName.java 的源码解析逻辑中,这一变化被直接体现出来:解析器通过按/切分(限制为 4 段),当结果恰好为 3 段时按新格式(tenant/namespace/localName)解析并将cluster置为null;当结果为 4 段时按旧格式(tenant/cluster/namespace/localName)解析。
兼容性说明:使用旧版名称格式的既有 Topic 将继续正常工作,无需任何修改,Pulsar 也没有计划取消这种兼容性。从源码也可以印证:
TopicName同时保留了对新旧两种格式的解析分支,并将 legacy 名称归一化处理。
灵活命名(Flexible topic naming)
虽然 Pulsar 2.0 中所有 Topic 名称在内部都采用{persistent|non-persistent}://tenant/namespace/topic的完整形式,但出于简化使用的考虑,大多数场景下现在可以使用简写名称。这一灵活命名体系的根源在于,Pulsar 2.0 引入了默认的 Topic 类型、租户和命名空间:
| Topic 维度 | 默认值 |
|---|---|
| Topic 类型(topic type) | persistent |
| 租户(tenant) | public |
| 命名空间(namespace) | default |
也就是说,当你省略某一部分时,Pulsar 会将其自动补齐为默认值。下表给出了一些简写名称到完整名称的转换示例:
| 输入的 Topic 名称 | 转换后的完整 Topic 名称 |
|---|---|
my-topic | persistent://public/default/my-topic |
my-tenant/my-namespace/my-topic | persistent://my-tenant/my-namespace/my-topic |
这一默认值逻辑在 TopicName.java 中同样有明确的常量定义:PUBLIC_TENANT = "public"、DEFAULT_NAMESPACE = "default"。当传入的名称不含://(即短名称)时,解析器会按如下规则补全(见 TopicName.java):
- 名称形如
tenant/namespace/topic(3 段):自动添加persistent://前缀; - 名称只有一个单词(1 段,即
topic):自动补全为persistent://public/default/topic; - 其他形式(如 2 段):抛出
IllegalArgumentException提示应使用<tenant>/<namespace>/<topic>或<topic>格式。
非持久化 Topic 的例外情况
需要注意的是,上述基于默认值的简写规则仅适用于持久化(persistent)Topic。对于非持久化 Topic,你必须指定完整的 Topic 名称。
例如,你不能使用类似non-persistent://my-topic的简写名称,而必须写成:
non-persistent://public/default/my-topic从 concepts-messaging.md 可以看到,非持久化 Topic 的完整形式为non-persistent://tenant/namespace/topic,生产者与消费者可以像连接持久化 Topic 一样连接非持久化 Topic,唯一的区别是名称必须以non-persistent开头,且支持 exclusive、shared、failover 三种订阅类型。
无需显式创建 Topic
与 Topic 名称相关的另一个重要特性是:Pulsar 中无需显式创建 Topic。如果客户端尝试向一个尚不存在的 Topic 写入或接收消息,Pulsar 会自动在 Topic 名称指定的 namespace 下创建该 Topic;如果客户端创建 Topic 时未指定 tenant 和 namespace,Topic 将创建在默认的public租户与default命名空间中(见 concepts-messaging.md)。
总结:迁移到 Pulsar 2.0+ 命名体系的关键清单
Pulsar 2.0 的命名体系变革可以从三个层面把握:
- 术语层面:
property全面更名为tenant,命令行工具由pulsar-admin properties迁移至pulsar-admin tenants,旧术语仅保留为废弃兼容; - Topic 名称结构层面:由
{persistent|non-persistent}://property/cluster/namespace/topic简化为{persistent|non-persistent}://tenant/namespace/topic,cluster 组件被移除,旧格式 Topic 继续兼容运行; - 日常使用层面:持久化 Topic 支持基于默认值(
persistent、public、default)的简写命名,例如my-topic即代表persistent://public/default/my-topic;而非持久化 Topic 必须始终写全名,且 Topic 本地名称中不允许出现/。
配合 Pulsar 2.0 新增的 Pulsar Functions 轻量计算能力(其部署命令中的--tenant/--namespace参数正是新命名体系的直接应用),你可以用更简洁、更清晰的方式组织和管理 Pulsar 集群中的多租户资源与消息流。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 2.0 升级指南:简化 Topic 命名、Tenant 术语重构与 Pulsar Functions 新特性解读
Apache Pulsar 2.0 升级指南:简化 Topic 命名、Tenant 术语重构与 Pulsar Functions 新特性解读 Pulsar 2.
消息队列后端流处理Apache Pulsar 多租户机制详解:Tenant、Namespace 与基于 `__change_events` 系统主题的 Topic 级策略
Apache Pulsar 多租户机制详解:Tenant、Namespace 与基于 __change_events 系统主题的 Topic 级策略 Apach
消息队列后端流处理Apache Pulsar Flume Source Connector 实战指南:将 Flume Agent 日志导入 Pulsar Topic
Apache Pulsar Flume Source Connector 实战指南:将 Flume Agent 日志导入 Pulsar Topic Flume
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考