☰
Apache Pulsar Geo-Replication 跨集群复制原理与配置实战指南
2026/9/28 20:37:17 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

本文基于 Apache Pulsar 官方文档,系统讲解 Pulsar 的多集群复制(Geo-replication)机制:从复制的基本工作原理、租户(Tenant)级权限模型、命名空间级复制配置,到选择性复制、Topic 统计与垃圾回收等实战要点,并结合本仓库源码与配置文件,帮助读者掌握在多数据中心场景下搭建跨地域消息多活的具体方法。

什么是 Geo-replication

Geo-replication(地理复制)是 Pulsar 实例(Pulsar Instance)中多个集群之间对持久化存储消息数据的复制机制。它允许同一个主题(Topic)的消息在多个 Pulsar 集群中同时存在并被各自集群的消费者独立消费,是构建跨地域(多活)消息系统的核心能力。

Pulsar 的多集群(Multi-cluster)部署模型下,一个 Pulsar 实例可以包含任意数量的 Pulsar 集群。默认情况下,消息只保存在发布它的那个集群中;而启用 Geo-replication 后,消息会被自动复制到其他配置了复制关系的集群中,从而实现多区域的数据冗余与就近消费。

工作原理图

下图展示了 Pulsar 跨集群复制的整体过程(该图来自仓库 site2/website-next/static/assets/geo-replication.png):

在上图中,P1、P2、P3三个生产者分别向Cluster-A、Cluster-B、Cluster-C三个集群上的T1主题发布消息,这些消息会即时地被复制到所有参与复制的集群。复制完成后,C1与C2消费者就可以从各自所在的集群中消费到全部三处发布的消息。

如果没有启用 Geo-replication,C1和C2消费者将无法消费由P3生产者(位于其他集群)发布的消息。换句话说,Geo-replication 让任意集群中的生产者与消费者都能跨越地域边界,共享同一份逻辑上的消息数据。

Geo-replication 与 Pulsar Properties(租户)

Pulsar 中 Geo-replication 的启用是**按租户(per-tenant)**进行的。只有在创建一个同时允许访问两个集群的租户之后,这两个集群之间才能启用复制。

尽管复制关系最终发生在两个集群之间,但它的实际管理粒度是**命名空间(Namespace)**级别。要为一个命名空间启用 Geo-replication,需要完成以下两步:

  1. 启用复制命名空间(详见下文 Enabling geo-replication namespaces);
  2. 将该命名空间配置为跨两个或更多已配给(provisioned)集群复制。

一旦完成配置,该命名空间下任意主题上发布的消息,都会被复制到配置集合中的所有集群。

关于“Properties”术语的说明:Pulsar 早期的术语体系中,“Property”即现在的“Tenant”(租户)。因此本文档中的“per-property/tenant 级别”指的是在租户层面控制集群访问权限,而复制配置本身落在命名空间层面。相关术语可参考仓库文档 reference-terminology.md。

本地持久化与转发机制

当一个消息被发布到 Pulsar 主题时,它首先在本地集群中完成持久化,然后异步转发到远端集群。具体来说:

  • 在正常情况下(无网络问题),消息会在与分发到本地消费者同时被立即复制出去;
  • 端到端的投递延迟通常取决于远端区域之间的网络往返时间(RTT);
  • 即使远端集群暂时不可达(例如发生网络分区),应用也仍然可以在任何一个集群中创建生产者和消费者——本地生产与消费不受影响,消息会积压在复制通道中,待网络恢复后继续转发。

订阅(Subscription)是集群本地的

需要注意一个关键特性:订阅是集群本地的(local to a cluster)。虽然生产者和消费者可以在 Pulsar 实例中的任何集群发布与消费,但订阅只在创建它的集群内生效,且不能跨集群迁移。如果需要迁移订阅,必须在目标集群中新建一个订阅。

Subscriptions are local to a cluster

在复制场景中,这一点意味着:C1和C2消费者各自消费的是本集群内订阅所跟踪的消息进度,不同集群的消费进度彼此独立。

在上文的三集群示例中,T1主题在Cluster-A、Cluster-B、Cluster-C之间复制。三个集群中任何集群产生的消息,都会被投递到其他集群中的所有订阅。因此:

  • C1、C2消费者会收到P1、P2、P3三个生产者发布的所有消息;
  • 消息的顺序性在单生产者(per-producer)维度依然得到保证。

源码层面的实现佐证

从源码结构看,复制通道由 broker 内部的 Replicator 组件实现,核心类位于 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java 及其持久化子类 PersistentReplicator.java(非持久化主题对应 NonPersistentReplicator.java)。可以推断出复制通道的运行方式:

  • 每个复制通道在本地主题上维护一个复制游标(cursor)(见PersistentReplicator中的ManagedCursor cursor),从该游标位置持续读取本地已持久化的消息条目(entries);
  • 复制通道通过一个内部创建的**复制生产者(replication producer)**将消息转发到远端集群的同一主题,其生产者命名格式为本地集群-->远端集群(分隔符常量为REPL_PRODUCER_NAME_DELIMITER = "-->",见 AbstractReplicator.java);
  • 复制生产者禁用了消息批处理(enableBatching(false))、设置了独立的发送超时与待处理消息队列上限(maxPendingMessages(producerQueueSize)),从而保证异步转发过程中出现网络故障时消息会在本地游标处积压、恢复后继续复制,这正是“本地持久化 + 异步转发”的实现基础。

配置复制(Configuring Replication)

如前述,Pulsar 的 Geo-replication 是在租户(tenant)层面管理的。下面按顺序给出完整的配置步骤。

向租户授予集群访问权限

要将消息复制到某个集群,租户必须拥有使用该集群的权限。你可以在创建租户时一并授予,也可以事后补授。

在创建租户时指定所有目标集群:

$ bin/pulsar-admin tenants create my-tenant \ --admin-roles my-admin-role \ --allowed-clusters us-west,us-east,us-cent

其中:

  • --admin-roles:指定该租户的管理员角色;
  • --allowed-clusters:指定该租户允许访问的集群列表,只有列出的集群才能参与该租户命名空间的复制。

为已有租户更新权限:将create换成update即可,例如:

$ bin/pulsar-admin tenants update my-tenant \ --admin-roles my-admin-role \ --allowed-clusters us-west,us-east,us-cent

启用复制命名空间(Enabling geo-replication namespaces)

创建命名空间的命令如下:

$ bin/pulsar-admin namespaces create my-tenant/my-namespace

创建之初,命名空间未绑定任何集群。需要通过set-clusters子命令将命名空间指派给多个集群:

$ bin/pulsar-admin namespaces set-clusters my-tenant/my-namespace \ --clusters us-west,us-east,us-cent

命名空间的复制集群集合可以随时修改,且不会中断正在进行的流量:一旦配置变更,所有集群中的复制通道会立即建立或停止。从源码结构看,set-clusters的底层实现会更新命名空间的复制策略,并触发 broker 对复制通道的动态增删(相关实现可参见 pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java 中关于复制集群的校验与更新逻辑)。

在 Geo-replication 主题上使用 Topic

一旦创建了复制命名空间,生产者或消费者在该命名空间内创建的任何主题都会自动跨集群复制。通常情况下,每个应用只需使用本地集群的serviceUrl连接即可——无论是生产还是消费,都走本地集群,由 broker 的复制通道负责跨集群数据同步。

选择性复制(Selective replication)

默认情况下,消息会被复制到该命名空间配置的所有集群。你可以通过为单条消息指定**复制列表(replication list)**来限制其复制范围:只有列表中的集群会收到这条消息。

以下为Java API示例(注意构造 {@inject: javadoc:Message:/client/org/apache/pulsar/client/api/Message} 对象时使用的setReplicationClusters方法):

List<String> restrictReplicationTo = Arrays.asList( "us-west", "us-east" ); Producer producer = client.newProducer() .topic("some-topic") .create(); producer.newMessage() .value("my-payload".getBytes()) .setReplicationClusters(restrictReplicationTo) .send();

说明:

  • setReplicationClusters(List<String>)定义于客户端 API 的TypedMessageBuilder接口(见 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/TypedMessageBuilder.java),用于设置单条消息的复制集群白名单;
  • 该字段在消息元数据中随消息传递,broker 侧的复制通道会读取并据此决定是否将该消息转发到对应集群;
  • 上例中消息只会复制到us-west与us-east两个集群,即使命名空间还配置了us-cent也不会复制过去。
Topic 统计(Topic stats)

Geo-replication 主题的主题级统计信息可以通过pulsar-admin命令行工具以及 REST API 获取:

$ bin/pulsar-admin persistent stats persistent://my-tenant/my-namespace/my-topic

每个集群只上报本地的统计信息,其中包括:

  • 复制消息的入站速率(incoming replication rate)与出站速率(outgoing replication rate);
  • 复制通道的积压量(backlog)。

从源码看,broker 侧使用ReplicatorStatsImpl(位于 pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/ReplicatorStatsImpl.java)记录每个复制通道的发送速率、接收速率与积压等指标,PersistentReplicator中也维护了msgOut、msgExpired等速率统计对象(见 PersistentReplicator.java)。

删除一个 Geo-replication 主题

由于 Geo-replication 主题同时存在于多个区域,无法直接删除某个复制主题。正确的做法是依赖**自动主题垃圾回收(automatic topic garbage collection)**机制。

在 Pulsar 中,一个主题在同时满足以下三个条件时会被自动删除:

  1. 没有生产者或消费者连接到该主题;
  2. 该主题上没有订阅(subscriptions);
  3. 没有因保留策略(retention)而继续保留的消息。

对于 Geo-replication 主题,每个区域都会使用一种容错机制来独立判断何时可以安全地在本地删除该主题。

关闭垃圾回收:你可以通过将 broker 配置中的brokerDeleteInactiveTopicsEnabled设置为false来显式禁用主题垃圾回收(相关 broker 配置项说明见 reference-configuration.md)。

删除一个 Geo-replication 主题的正确步骤:

  1. 关闭该主题上的所有生产者和消费者;
  2. 在每一个复制集群中删除该主题的所有本地订阅;
  3. 当 Pulsar 判定整个系统中该主题不再存在任何有效订阅时,就会对主题进行垃圾回收。

与复制相关的 Broker 配置

Geo-replication 的运行时行为还受 broker 配置影响。仓库中的 conf/broker.conf 提供了相关默认值,例如:

配置项默认值含义
replicationConnectionsPerBroker16每台 broker 用于复制连接的连接数上限
replicationProducerQueueSize1000复制生产者的待处理消息队列大小,即复制通道允许积压的未确认消息数
replicationPolicyCheckDurationSeconds600周期性检查复制策略的间隔(秒),用于避免复制器(replicator)出现不一致状态
replicationMetricsEnabledtrue是否启用复制相关指标(metrics)的采集
replicationTlsEnabledfalse复制连接是否启用 TLS

其中replicationProducerQueueSize直接对应AbstractReplicator构造时对内部复制生产者设置的maxPendingMessages(producerQueueSize)(见 AbstractReplicator.java):该值决定了复制通道在网络故障期间能够在本地积压的待发送消息上限,是调优跨区域复制吞吐与背压的关键参数之一。

另外,conf/broker.conf 中还有复制分发限流配置(默认值为0,表示不限流),可用于控制每个复制器(replicator)的消息/字节分发速率,防止复制流量挤占本地消费带宽。

常见场景与注意事项

  • 就近读写:复制命名空间下的应用始终连接本地集群的serviceUrl,生产与消费都发生在本地,跨集群的数据同步由 broker 后台完成,因此延迟主要受本地集群与远端集群之间的 RTT 影响。
  • 多活消费:所有集群中的消费者都能独立消费到全量消息,但每个集群的订阅进度互不相同;请勿期望不同集群的消费者共享同一个消费进度。
  • 复制范围的动态调整:命名空间的复制集群集合可以随时通过set-clusters调整,复制通道会即时增删,适用于扩缩容或灾备演练场景。
  • 删除主题的约束:复制主题不能直接删除,必须先在全部复制集群中清理订阅、关闭连接,等待自动垃圾回收完成。
  • 配置参考:完整的租户、命名空间与集群管理命令可参见 reference-pulsar-admin.md;broker 全部复制相关配置项见 reference-configuration.md 与 conf/broker.conf。

小结

Apache Pulsar 的 Geo-replication 以“租户授权 + 命名空间绑定集群”为配置模型,以“本地持久化、异步转发”为运行机制,配合选择性复制、复制统计与自动垃圾回收等能力,为多数据中心、跨地域多活的消息系统提供了开箱即用的官方方案。本文梳理的配置命令与源码实现,可直接用于在真实集群上搭建跨区域复制链路。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

相关推荐

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

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

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

立即咨询