集群与负载均衡
集群
Camel 提供了以下与集群相关的 SPI:
集群服务(Cluster Service)
一个常规的 Camel 服务,用于管理集群资源,例如视图(见下文)
集群视图(Cluster View)
表示集群的一个视图,拥有自己独立的一组资源。目前视图提供以下支持:
- 领导者选举
- 拓扑事件,例如成员加入/离开集群
集群成员(Cluster Member)
表示集群中的一个成员。
集群 SPI 配置
集群服务与其他 Camel 服务并无二致,因此要对其进行配置,只需将你的实现注册到 CamelContext 中:
仅 Java:在 CamelContext 上注册集群服务
MyClusterServiceImpl service = new MyClusterServiceImpl();
context.addService(service);Cluster Service(集群服务)的配置取决于你所选择的实现。Camel 开箱即提供了以下几种实现:
| 类型 | 模块 | 类 |
|---|---|---|
| consul | camel-consul | org.apache.camel.component.consul.cluster.ConsulClusterService |
| file | camel-file | org.apache.camel.component.file.cluster.FileLockClusterService |
| infinispan | camel-infinispan | org.apache.camel.component.infinispan.cluster.InfinispanClusterService |
| jgroups-raft | camel-jgroups-raft | org.apache.camel.component.jgroups.raft.cluster.JGroupsRaftClusterService |
| kubernetes | camel-kubernetes | org.apache.camel.component.kubernetes.cluster.KubernetesClusterService |
| zookeeper | camel-zookeeper | org.apache.camel.component.zookeeper.cluster.ZooKeeperClusterService |
配置选项:
ConsulClusterService
| 名称 | 说明 | 默认值 | 类型 |
|---|---|---|---|
| sessionTtl | Consul 会话的 TTL(秒) | 60 | int |
| sessionLockDelay | Consul 会话的锁延迟时间(秒) | 5 | int |
| sessionRefreshInterval | Consul 会话的刷新间隔(秒) | 5 | int |
| rootPath | Consul 集群的根目录路径 | /camel | String |
FileLockClusterService
| 名称 | 说明 | 默认值 | 类型 |
|---|---|---|---|
| acquireLockDelay | 在开始尝试获取集群锁之前需要等待的时间。注意,如果 FileLockClusterService 判断没有集群成员在运行,或者无法可靠地确定集群状态,则初始延迟将根据 acquireLockDelayInterval 计算得出 | 1 | long |
| acquireLockDelayUnit | acquireLockDelay 的时间单位 | SECONDS | TimeUnit |
| acquireLockInterval | 两次尝试获取集群锁之间等待的时间,按挂钟时间计算。所有集群成员必须使用相同的值,以确保领导权检查和领导者存活检测保持一致 | 10 | long |
| acquireLockIntervalUnit | acquireLockInterval 的时间单位 | SECONDS | TimeUnit |
| clusterDataTaskMaxAttempts | 设置集群数据任务的运行次数,包括首次执行以及失败或超时后的后续重试次数。当集群数据根目录位于网络文件存储上、I/O 操作偶尔可能长时间或不可预测地阻塞时,此配置非常有用 | 5 | int |
| clusterDataTaskTimeout | 设置集群数据任务(读取或写入集群数据)的超时时间。当集群数据根目录位于网络存储上、I/O 操作偶尔可能长时间或不可预测地阻塞时,超时机制非常有用 | 10 | long |
| clusterDataTaskTimeoutUnit | clusterDataTaskTimeoutUnit 的时间单位 | SECONDS | TimeUnit |
| heartbeatTimeoutMultiplier | 应用于集群领导者 acquireLockDelayInterval 的乘数,用于确定跟随者在认为领导者"过期"之前应等待多长时间。例如,如果领导者每 2 秒更新一次心跳且 heartbeatTimeoutMultiplier 为 3,则跟随者最多可容忍 2s * 3 = 6s 的静默时间,超过后将宣布领导者不可用 | 5 | int |
| 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 | 注册到清单中的本地集群成员对应的缓存条目的存活时间 | 30 | long |
| lifespanTimeUnit | 存活时间的时间单位 | SECONDS | TimeUnit |
JGroupsRaftClusterService
| 名称 | 描述 | 默认值 | 类型 |
|---|---|---|---|
| jgroupsConfig | JGroups Raft 配置文件的路径 | raft.xml | String |
| jgroupsClusterName | 集群的名称 | jgroupsraft-master | String |
| raftHandle | RaftHandle 实例 | org.jgroups.raft.RaftHandle | |
| raftId | 唯一的 Raft 标识 | String |
KubernetesClusterService
| 名称 | 描述 | 默认值 | 类型 |
|---|---|---|---|
| leaseResourceType | 用于保存租约的 Kubernetes 资源类型 | LeaseResourceType.Lease | LeaseResourceType |
| kubernetesResourcesNamespace | 包含 Pod 和用于加锁的 ConfigMap 的 Kubernetes 命名空间 | String | |
| kubernetesResourceName | 用于加锁的资源名称(若使用多个资源则为前缀) | leaders | String |
| groupName | 所选 ConfigMap 中的锁组名称(根据 Camel 集群约定,即命名空间) | String | |
| podName | 当前 Pod 的名称(默认为主机名) | String | |
| clusterLabels | 用于标识集群成员的标签 | 空映射 | Map |
| jitterFactor | 用于避免所有 Pod 在同一时刻调用 Kubernetes API 的抖动因子 | 1.2 | double |
| leaseDurationMillis | 当前领导者租约的默认持续时间 | 15000 | long |
| renewDeadlineMillis | 超过该期限后,领导者必须停止其服务,因为它可能已经失去领导地位 | 10000 | long |
| retryPeriodMillis | 两次检查并获取领导地位的尝试之间的时间间隔。该时间会使用抖动因子进行随机化 | 2000 | long |
ZooKeeperClusterService
| 名称 | 描述 | 默认值 | 类型 |
|---|---|---|---|
| nodes | Zookeeper 服务器主机(多个服务器可用逗号分隔) | List | |
| namespace | ZooKeeper 命名空间。如果在此处设置了命名空间,所有路径都会加上该命名空间前缀 | String | |
| reconnectBaseSleepTime | 重试之间初始等待时间 | long | |
| reconnectBaseSleepTimeUnit | ReconnectBaseSleepTime 的时间单位。默认为 | MILLISECONDS | TimeUnit |
| reconnectMaxRetries | 最大重试次数 | 3 | int |
| sessionTimeout | 会话超时时间(毫秒) | 60000 | long |
| sessionTimeoutUnit | 会话超时的时间单位 | MILLISECONDS | TimeUnit |
| connectionTimeout | 连接超时时间(毫秒) | 15000 | long |
| connectionTimeoutUnit | 连接超时的时间单位 | TimeUnit.MILLISECONDS | TimeUnit |
| authInfoList | 包含 scheme 和 auth 的 AuthInfo 对象列表 | List | |
| maxCloseWait | 关闭期间等待后台线程结束的时间 | 1000 | long |
| maxCloseWaitUnit | MaxCloseWait 的时间单位 | MILLISECONDS | TimeUnit |
| 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 = falseMaster 组件
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) ...");
}
};
}评论
登录后参与评论
KnowForge