架构

集群与负载均衡

师成师成· 更新于 2026-09-28· 阅读 23 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

集群

Camel 提供了以下与集群相关的 SPI:

  • 集群服务(Cluster Service)

    一个常规的 Camel 服务,用于管理集群资源,例如视图(见下文)

  • 集群视图(Cluster View)

    表示集群的一个视图,拥有自己独立的一组资源。目前视图提供以下支持:

    • 领导者选举
    • 拓扑事件,例如成员加入/离开集群
  • 集群成员(Cluster Member)

    表示集群中的一个成员。

集群 SPI 配置

集群服务与其他 Camel 服务并无二致,因此要对其进行配置,只需将你的实现注册到 CamelContext 中:

仅 Java:在 CamelContext 上注册集群服务

MyClusterServiceImpl service = new MyClusterServiceImpl();
context.addService(service);

Cluster Service(集群服务)的配置取决于你所选择的实现。Camel 开箱即提供了以下几种实现:

类型模块类
consulcamel-consulorg.apache.camel.component.consul.cluster.ConsulClusterService
filecamel-fileorg.apache.camel.component.file.cluster.FileLockClusterService
infinispancamel-infinispanorg.apache.camel.component.infinispan.cluster.InfinispanClusterService
jgroups-raftcamel-jgroups-raftorg.apache.camel.component.jgroups.raft.cluster.JGroupsRaftClusterService
kubernetescamel-kubernetesorg.apache.camel.component.kubernetes.cluster.KubernetesClusterService
zookeepercamel-zookeeperorg.apache.camel.component.zookeeper.cluster.ZooKeeperClusterService

配置选项:

ConsulClusterService

名称说明默认值类型
sessionTtlConsul 会话的 TTL(秒)60int
sessionLockDelayConsul 会话的锁延迟时间(秒)5int
sessionRefreshIntervalConsul 会话的刷新间隔(秒)5int
rootPathConsul 集群的根目录路径/camelString

FileLockClusterService

名称说明默认值类型
acquireLockDelay在开始尝试获取集群锁之前需要等待的时间。注意,如果 FileLockClusterService 判断没有集群成员在运行,或者无法可靠地确定集群状态,则初始延迟将根据 acquireLockDelayInterval 计算得出1long
acquireLockDelayUnitacquireLockDelay 的时间单位SECONDSTimeUnit
acquireLockInterval两次尝试获取集群锁之间等待的时间,按挂钟时间计算。所有集群成员必须使用相同的值,以确保领导权检查和领导者存活检测保持一致10long
acquireLockIntervalUnitacquireLockInterval 的时间单位SECONDSTimeUnit
clusterDataTaskMaxAttempts设置集群数据任务的运行次数,包括首次执行以及失败或超时后的后续重试次数。当集群数据根目录位于网络文件存储上、I/O 操作偶尔可能长时间或不可预测地阻塞时,此配置非常有用5int
clusterDataTaskTimeout设置集群数据任务(读取或写入集群数据)的超时时间。当集群数据根目录位于网络存储上、I/O 操作偶尔可能长时间或不可预测地阻塞时,超时机制非常有用10long
clusterDataTaskTimeoutUnitclusterDataTaskTimeoutUnit 的时间单位SECONDSTimeUnit
heartbeatTimeoutMultiplier应用于集群领导者 acquireLockDelayInterval 的乘数,用于确定跟随者在认为领导者"过期"之前应等待多长时间。例如,如果领导者每 2 秒更新一次心跳且 heartbeatTimeoutMultiplier 为 3,则跟随者最多可容忍 2s * 3 = 6s 的静默时间,超过后将宣布领导者不可用5int
rootPath文件集群根目录路径String
与大多数集群实现一样,文件锁集群中所有机器的时钟保持同步非常重要。

将集群根目录存储在 NFS 上

当将 FileLockClusterService 的集群根目录(参见 rootPath 配置选项)存储在 NFS 上时,所有集群成员必须使用允许故障切换可靠工作的 NFS 客户端挂载选项。

以下 NFS 客户端挂载选项有助于实现这一点。你可以根据实际需要调整这些选项。

  • soft — 当 NFS 服务器未响应请求时,禁止客户端持续进行重传尝试。
  • timeo=10 — NFS 客户端在再次发送请求之前,等待 NFS 服务器响应的时间,单位为十分之一秒。默认值为 600(60 秒)。
  • retrans=1 — 指定 NFS 客户端尝试向 NFS 服务器重传失败请求的次数。默认值为 3。客户端在每次 retrans 尝试之间会等待一个 timeo 超时周期。
  • lookupcache=none — 指定内核如何管理挂载点的目录项缓存。none 会强制客户端在使用所有缓存项之前对其进行重新验证。这使集群领导者能够立即检测到对其锁文件所做的任何更改,并防止锁检查机制返回不正确的有效性信息。
  • sync — 任何向挂载点上的文件写入数据的系统调用,都会在将控制权返回给用户空间之前先将数据刷新到 NFS 服务器。该选项可提供更高的数据缓存一致性。

有关 NFS 挂载选项的更多信息,请参阅 http://linux.die.net/man/5/nfs。

InfinispanClusterService

名称说明默认值类型
lifespan注册到清单中的本地集群成员对应的缓存条目的存活时间30long
lifespanTimeUnit存活时间的时间单位SECONDSTimeUnit

JGroupsRaftClusterService

名称描述默认值类型
jgroupsConfigJGroups Raft 配置文件的路径raft.xmlString
jgroupsClusterName集群的名称jgroupsraft-masterString
raftHandleRaftHandle 实例org.jgroups.raft.RaftHandle
raftId唯一的 Raft 标识String

KubernetesClusterService

名称描述默认值类型
leaseResourceType用于保存租约的 Kubernetes 资源类型LeaseResourceType.LeaseLeaseResourceType
kubernetesResourcesNamespace包含 Pod 和用于加锁的 ConfigMap 的 Kubernetes 命名空间String
kubernetesResourceName用于加锁的资源名称(若使用多个资源则为前缀)leadersString
groupName所选 ConfigMap 中的锁组名称(根据 Camel 集群约定,即命名空间)String
podName当前 Pod 的名称(默认为主机名)String
clusterLabels用于标识集群成员的标签空映射Map
jitterFactor用于避免所有 Pod 在同一时刻调用 Kubernetes API 的抖动因子1.2double
leaseDurationMillis当前领导者租约的默认持续时间15000long
renewDeadlineMillis超过该期限后,领导者必须停止其服务,因为它可能已经失去领导地位10000long
retryPeriodMillis两次检查并获取领导地位的尝试之间的时间间隔。该时间会使用抖动因子进行随机化2000long

ZooKeeperClusterService

名称描述默认值类型
nodesZookeeper 服务器主机(多个服务器可用逗号分隔)List
namespaceZooKeeper 命名空间。如果在此处设置了命名空间,所有路径都会加上该命名空间前缀String
reconnectBaseSleepTime重试之间初始等待时间long
reconnectBaseSleepTimeUnitReconnectBaseSleepTime 的时间单位。默认为MILLISECONDSTimeUnit
reconnectMaxRetries最大重试次数3int
sessionTimeout会话超时时间(毫秒)60000long
sessionTimeoutUnit会话超时的时间单位MILLISECONDSTimeUnit
connectionTimeout连接超时时间(毫秒)15000long
connectionTimeoutUnit连接超时的时间单位TimeUnit.MILLISECONDSTimeUnit
authInfoList包含 scheme 和 auth 的 AuthInfo 对象列表List
maxCloseWait关闭期间等待后台线程结束的时间1000long
maxCloseWaitUnitMaxCloseWait 的时间单位MILLISECONDSTimeUnit
retryPolicy要使用的重试策略。RetryPolicy
basePath在 ZooKeeper 中存储的基础路径String

配置示例:

  • Spring Boot

    camel.cluster.file.enabled = true
    camel.cluster.file.id = ${random.uuid}
    camel.cluster.file.root = ${java.io.tmpdir}
  • Spring XML

    <?xml version="1.0" encoding="UTF-8"?>
    <beans xmlns="http://www.springframework.org/schema/beans"
           xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
           xsi:schemaLocation="
             http://www.springframework.org/schema/beans
             http://www.springframework.org/schema/beans/spring-beans.xsd
             http://camel.apache.org/schema/spring
             http://camel.apache.org/schema/spring/camel-spring.xsd">
    
      <bean id="zx" class="org.apache.camel.component.zookeeper.cluster.ZooKeeperClusterService">
        <property name="id" value="node-1"/>
        <property name="basePath" value="/camel/cluster"/>
        <property name="nodes" value="localhost:2181"/>
      </bean>
    
      <camelContext xmlns="http://camel.apache.org/schema/spring" autoStartup="false">
        ...
      </camelContext>
    
    </beans>

Cluster SPI 的使用

Cluster SPI 被以下新增实现所使用:

  • ClusteredRoutePolicy

    这是 RoutePolicy 的一个实现,当它所使用的集群视图(Cluster View)获得领导权时,会启动与其关联的路由。

    仅 Java:在路由中使用 ClusteredRoutePolicy

    context.addRoutes(new RouteBuilder {
        @Override
        public void configure() throws Exception {
            // Create the route policy
            RoutePolicy policy = ClusteredRoutePolicy.forNamespace("my-ns");
    
            // bind the policy to one or more routes
            from("timer:clustered?delay=1000&period=1000")
                .routePolicy(policy)
                .log("Route ${routeId} is running ...");
        }
    });

要将同一策略应用于所有路由,可以使用专用的 RoutePolicyFactory。

仅 Java:使用 ClusteredRoutePolicyFactory

// 将集群路由策略工厂添加到上下文
context.addRoutePolicyFactory(ClusteredRoutePolicyFactory.forNamespace("my-ns"));

context.addRoutes(new RouteBuilder {
    @Override
    public void configure() throws Exception {
        // 将策略绑定到一个或多个路由
        from("timer:clustered?delay=1000&period=1000")
            .log("Route ${routeId} is running ...");
    }
});
  • ClusteredRouteController

    这是 RouteController SPI 的一种实现,它让 camel context 先启动,然后在取得/失去领导权时启动/停止路由。该实现与 spring-boot 应用集成良好,假设你的路由配置如下:

    仅 Java:ClusteredRouteController 的 Spring Boot 路由定义

    @Bean
    public RouteBuilder routeBuilder() {
        return new RouteBuilder() {
            @Override
            public void configure() throws Exception {
                from("timer:heartbeat?period=10000")
                    .routeId("heartbeat")
                    .log("HeartBeat route (timer) ...");
                from("timer:clustered?period=5000")
                    .routeId("clustered")
                    .log("Clustered route (timer) ...");
            }
        };
    }

然后,你可以利用 Spring Boot 配置让它们组成集群:

# 启用路由控制器
camel.clustered.controller.enabled = true

# 定义路由的默认命名空间
camel.clustered.controller.namespace = my-ns

# 将 id 为 'heartbeat' 的路由排除在集群之外
camel.clustered.controller.routes[heartbeat].clustered = false
  • Master 组件

    master 组件与 ClusteredRoutePolicy 类似,但它工作在消费者级别,因此能确保在任何时刻集群中只有一个端点在消费资源。它的配置非常简单,你只需按照 master 组件的语法为单例端点加上前缀即可:

    master:namespace:delegateUri

具体示例:

仅限 Java:使用 master 组件的 Spring Boot 路由

@Bean
public RouteBuilder routeBuilder() {
    return new RouteBuilder() {
        @Override
        public void configure() throws Exception {
            from("timer:heartbeat?period=10000")
                .routeId("heartbeat")
                .log("HeartBeat route (timer) ...");

            from("master:my-ns:timer:clustered?period=5000")
                .routeId("clustered")
                .log("Clustered route (timer) ...");
        }
    };
}

评论

登录后参与评论

正在加载评论…