大家好,本文给大家介绍一下Elastic-Job 中作业执行时选主节点,对作业进行分片的逻辑
选主节点执行分片逻辑的过程
文 | 宋小生
7.4.5 分片逻辑
先来回顾下每个实例上的作业获取分片上下文的代码,在正常获取本作业实例被分配的分片列表之前,先进行判断是否需要重新分片处理:
@Override
public ShardingContexts getShardingContexts() {
boolean isFailover = configService.load(true).isFailover();
if (isFailover) {
List<Integer> failoverShardingItems = failoverService.getLocalFailoverItems();
if (!failoverShardingItems.isEmpty()) {
return executionContextService.getJobShardingContext(failoverShardingItems);
}
}
shardingService.shardingIfNecessary();
List<Integer> shardingItems = shardingService.getLocalShardingItems();
if (isFailover) {
shardingItems.removeAll(failoverService.getLocalTakeOffItems());
}
shardingItems.removeAll(executionService.getDisabledItems(shardingItems));
return executionContextService.getJobShardingContext(shardingItems);
}
接下来我们就看下elasticjob是如何进行分片的,调用shardingService类型下的shardingIfNecessary()方法来判断是否需要分片,如果需要分片则执行分片逻辑,具体代码如下:
public void shardingIfNecessary() {
//先获取当前可用实例
List<JobInstance> availableJobInstances = instanceService.getAvailableJobInstances();
//判断分片状态
if (!isNeedSharding() || availableJobInstances.isEmpty()) {
return;
}
//非主节点if判断结果为true,自旋等待主节点分片完毕,主节点进行分片,如果主节点不存在则先抢占主节点
if (!leaderService.isLeaderUntilBlock()) {
blockUntilShardingCompleted();
return;
}
//主节点:等待上次执行的作业执行完毕
waitingOtherJobCompleted();
//主节点:拉取作业最新分片配置信息,可用于实时生效
LiteJobConfiguration liteJobConfig = configService.load(false);
int shardingTotalCount = liteJobConfig.getTypeConfig().getCoreConfig().getShardingTotalCount();
log.debug("Job '{}' sharding begin.", jobName);
//主节点:分片开始持久化分片处理中状态机
jobNodeStorage.fillEphemeralJobNode(ShardingNode.PROCESSING, "");
//主节点:重置分片节点信息清理历史数据
resetShardingInfo(shardingTotalCount);
//获取分片策略算法
JobShardingStrategy jobShardingStrategy = JobShardingStrategyFactory.getStrategy(liteJobConfig.getJobShardingStrategyClass());
//主节点:使用Zookeeper事务进行分片
jobNodeStorage.executeInTransaction(new PersistShardingInfoTransactionExecutionCallback(jobShardingStrategy.sharding(availableJobInstances, jobName, shardingTotalCount)));
log.debug("Job '{}' sharding complete.", jobName);
}

图7.4 选主分片过程
详细来看下执行过程,
1)首先获取所有可用服务器实例和可用服务器实例需要满足如下条件:
作业实例临时节点存在: 代表进程存在,对应节点{作业名字}/instance/作业实例id
作业服务器持久节点存在: 并且值不是禁用状态,对应节点 {作业名字}/servers/作业ip ,禁用的状态值为DISABLED。
2)然后判断下分片状态和前面查询的可用实例是否满足条件不满足则返回。
是否需要分片,状态是以分片标示节(leader/sharding/necessary)存在来决定的,回顾下我们在哪些地方见到过设置了分片节点:
ShardingListenerManager监听器管理器中分片总数和服务器状态变化的时候则会创建需节点重新分片。
诊断服务中当前分片中存在未上线的实例时候写入重新分片标记。
初始化作业注册启动信息registerStartUpInfo的时候如果分片节点不存在情况下则会写入重新分片标记。
3) 然后就是选主分片过程
if (!leaderService.isLeaderUntilBlock()) {
blockUntilShardingCompleted();
return;
}
这里判断当前节点是否为主节点:
如果是主节点:跳出if,主节点执行后面的分片逻辑。
如果是非主节点:则进入if条件开始自旋等待主节点进行分片完毕,当主节点分片完毕时候,这里处于自旋的非主节点则直接返回跳出shardingIfNecessary方法。
7.4.5.1 首先我们要明白为什么要选主节点?
主节点存在的目的是为了有效的对作业进行分片,可以想一下如果在分配工作的时候没有一个领导,每个人都可以进行工作分配该如何分配,每个人都有自己的决策,无法形成一致的意见,这里的分片与分配工作类似,如果有多个机器同时分片,各自设置各自的分片状态,这样就很容易存在每个机器对作业下实例分片结果不同,这样的分配就比较混乱,无法决定按哪个实例的分片为准,所以选举一个主节点负责统一分片,将任务分配给每个节点,每个机器分到几个分片那就执行几次作业。
那接下来详细看下是触发选主的的逻辑:
先来看下阻塞判断当前节点是否为主节点的方法:
public boolean isLeaderUntilBlock() {
while (!hasLeader() && serverService.hasAvailableServers()) {
log.info("Leader is electing, waiting for {} ms", 100);
BlockUtils.waitingShortTime();
if (!JobRegistry.getInstance().isShutdown(jobName) && serverService.isAvailableServer(JobRegistry.getInstance().getJobInstance(jobName).getIp())) {
electLeader();
}
}
return isLeader();
}
这里是个循环判断当前实例是否为主节点,如果不存在主节点并且有可用服务器,并且可用服务器是上线状态的话则休眠几秒,这个休眠是为了降低大量作业同时进入防止大量冲突,不过这个版本是休眠的固定时间100毫秒不太好后期的话官方在升级时候应该会去修改。
如果作业未关闭并且当前作业服务器可用则开始执行选举作业主节点的逻辑。
那如何选主呢,让我们来看下具体代码:
/**
* 选举主节点.
*/
public void electLeader() {
log.debug("Elect a new leader now.");
jobNodeStorage.executeInLeader(LeaderNode.LATCH, new LeaderElectionExecutionCallback());
log.debug("Leader election completed.");
}
leader节点被移除的时候,leader节点不存在,并且作业状态为启用状态的话也会触发选主。
/**
* 判断当前节点是否是主节点.
*
* @return 当前节点是否是主节点
*/
public boolean isLeader() {
return !JobRegistry.getInstance().isShutdown(jobName) && JobRegistry.getInstance().getJobInstance(jobName).getJobInstanceId().equals(jobNodeStorage.getJobNodeData(LeaderNode.INSTANCE));
}
判断方法比较简单主要看下当前节点是否是被主节点写入的实例,如果当前机器实例是主节点则返回true,否则返回false,再回过头到shardingIfNecessary分片方法里面如果当前实例不是作业对应主节点实例则执行。
选主节点完成了,接下来主节点开始进行分片逻辑,我们了解过选主的目的就是为了让主节点进行分片,有序的分配任务。如果当前不是主节点则需要的阻塞等待主节点分片完成后,再执行作业。主节点执行分片的逻辑等下在看,接下来我们先来了解下非主节点的阻塞等待分片完成的方法blockUntilShardingCompleted。
7.4.5.1 非主节点自旋等待主节点分片完成
blockUntilShardingCompleted阻塞分片完成再返回,作业后面作业执行的时候要获取当前实例对应的分片项拿到分片项才能执行作业,这里阻塞直到分片完成后才返回。
private void blockUntilShardingCompleted() {
while (!leaderService.isLeaderUntilBlock() && (jobNodeStorage.isJobNodeExisted(ShardingNode.NECESSARY) || jobNodeStorage.isJobNodeExisted(ShardingNode.PROCESSING))) {
log.debug("Job '{}' sleep short time until sharding completed.", jobName);
BlockUtils.waitingShortTime();
}
}
这里也是一个自旋(无限循环)判断当前不是主节点并且存在需要分片节点或者存在分片中节点的情况下则休眠等待,也就是说如果是非主节点则一直等待到分片完成(分片完成的情况下会删除分片相关的节点保证非主节点在这里的条件为false),这里等待的时候从节点主要通过分片的状态来判断是否继续等待,一般情况下我们发生了服务器状态变更后才会进行一次新的分片逻辑,后期分片存在节点则不会进入循环等待则直接获取分片执行作业。
非主分片进行了等待那看看主分片接下来该如何分片。
首先会处理幂等,当幂等配置开启的情况下不允许同一个分片同时多次执行的情况。这里如果幂等开关开启,有作业分片在运行的情况下不能进行分片防止出现同一个分片项序号在分片后被触发,导致幂等失效,这里要明白一点就是我们的幂等是针对同一个作业的同一个分片级别的。这里通过判断当前作业是否有存在的分片项正在执行,如果作业未结束仍旧在执行则需要等待运行中的作业分片完成之后才能分片。
private void waitingOtherJobCompleted() {
while (executionService.hasRunningItems()) {
log.debug("Job '{}' sleep short time until other job completed.", jobName);
BlockUtils.waitingShortTime();
}
}
public boolean hasRunningItems(final Collection<Integer> items) {
LiteJobConfiguration jobConfig = configService.load(true);
if (null == jobConfig || !jobConfig.isMonitorExecution()) {
return false;
}
for (int each : items) {
if (jobNodeStorage.isJobNodeExisted(ShardingNode.getRunningNode(each))) {
return true;
}
}
return false;
}
看一下这个过程的实现细节,这个方法是当前即将分片的主节点在分片之前如果配置开启了幂等执行配置monitorExecution则分片之前保证上次的所有节点作业全部执行完毕再去执行,如何保证呢,想一下如果有多台机器情况下我们如何知道其他机器上次的作业是否正在执行呢?,这个地方我们就要借助Zookeeper了,如果开启了幂等执行则在作业执行之前会在对应分片位置写入一个sharding/作业名字/running节点,作业执行完毕删除此节点,这个地方我们就可以判断所有节点是否存在这个运行中的节点,如果存在则等待上次所有作业执行完毕本次再去执行。
7.4.5.3 主节点重置分片节点
写入分片执行中临时节点leader/sharding/processing
jobNodeStorage.fillEphemeralJobNode(ShardingNode.PROCESSING, "");
//重置分片节点,重置分片节点就是将之前的分片节点删除,根据分片总数来重新创建分片项根节点{jobName}/sharding/{分片项目}
resetShardingInfo(shardingTotalCount);
private void resetShardingInfo(final int shardingTotalCount) {
for (int i = 0; i < shardingTotalCount; i++) {
jobNodeStorage.removeJobNodeIfExisted(ShardingNode.getInstanceNode(i));
jobNodeStorage.createJobNodeIfNeeded(ShardingNode.ROOT + "/" + i);
}
int actualShardingTotalCount = jobNodeStorage.getJobNodeChildrenKeys(ShardingNode.ROOT).size();
if (actualShardingTotalCount > shardingTotalCount) {
for (int i = shardingTotalCount; i < actualShardingTotalCount; i++) {
jobNodeStorage.removeJobNodeIfExisted(ShardingNode.ROOT + "/" + i);
}
}
}
重新设置分片信息就是先删除之前存在的分片,然后根据分片项重新创建新的分片节点,后期进行分片的时候只需要将对应节点信息写入分片项下即可。
7.4.5.4 主节点使用策略模式+Java反射来获取分片算法
读取分片策略,使用策略模式+Java反射根据不同分片算法策略类型来生成分片策略算法对象来执行不同的分片,使用策略模式+Java反射有效解决了当分片策略过多的时候产生的大量的判断。我们先来了解下实现过程。
JobShardingStrategy jobShardingStrategy = JobShardingStrategyFactory.getStrategy(liteJobConfig.getJobShardingStrategyClass());
public static JobShardingStrategy getStrategy(final String jobShardingStrategyClassName) {
if (Strings.isNullOrEmpty(jobShardingStrategyClassName)) {
return new AverageAllocationJobShardingStrategy();
}
try {
Class<?> jobShardingStrategyClass = Class.forName(jobShardingStrategyClassName);
if (!JobShardingStrategy.class.isAssignableFrom(jobShardingStrategyClass)) {
throw new JobConfigurationException("Class '%s' is not job strategy class", jobShardingStrategyClassName);
}
return (JobShardingStrategy) jobShardingStrategyClass.newInstance();
} catch (final ClassNotFoundException | InstantiationException | IllegalAccessException ex) {
throw new JobConfigurationException("Sharding strategy class '%s' config error, message details are '%s'", jobShardingStrategyClassName, ex.getMessage());
}
}
这里根据配置的分片策略类来在运行时生成分片策略对象,如果配置为空则使用平均分片策略类类型,否则的话判断一下配置的这个类是否满足是JobShardingStrategy类型的子类 如果是的话则在运行时生成对应策略对象。
分片的过程也就是把分片项分配到各个作业实例,有些作业实例可能分到0个有些作业实例可能会被分到多个分片,这个就要看具体的算法和分片总数了,接下来我们看下基于平均分配算法的分片策略:
常见分片算法配置有哪些呢:
AverageAllocationJobShardingStrategy: 基于平均分配算法的分片策略,也是默认的分片策略 。
OdevitySortByNameJobShardingStrategy :作业名的哈希值奇偶数决定IP升降序算法的分片策略。
RotateServerByNameJobShardingStrategy:作业名的哈希值对服务器列表进行轮转的分片策略。
7.4.5.5 主节点使用事务执行分片的逻辑
JobShardingStrategy jobShardingStrategy = JobShardingStrategyFactory.getStrategy(liteJobConfig.getJobShardingStrategyClass());
jobNodeStorage.executeInTransaction(new PersistShardingInfoTransactionExecutionCallback(jobShardingStrategy.sharding(availableJobInstances, jobName, shardingTotalCount)));
到这里主节点处理分片结束主节点分片结束后返回:shardingService的shardingIfNecessary()方法
非主节点则结束自旋跳出shardingIfNecessary()方法,分片进行完毕
接下来主节点和非主节点开始由当前实例信息获取当前实例被主节点分配的分片项列表,每个实例可以获取0个,1个或者多个分片项。
根据当前实例id 获取分片时候写入的分片项 :
获取当前机器对应的分片项
List<Integer> shardingItems = shardingService.getLocalShardingItems();
移除失效的分片项
if (isFailover) {
//如果是失效转移配置则移除这些分片项中失效的节点,这些失效的节点存在失效转移可执行节点,对应节点/作业名字/sharding/分片项/failover
shardingItems.removeAll(failoverService.getLocalTakeOffItems());
}
移除被禁用的分片项
//接下来就是移除被禁用的分片项 对应节点/作业名字/disabled/分片项目/disabled
shardingItems.removeAll(executionService.getDisabledItems(shardingItems));
封装分片上下文对象
//最后封装分片上下文对象
return executionContextService.getJobShardingContext(shardingItems);
这个分片上下文将会作为作业执行的参数接下来我们看看是分片上下文封装了哪些内容
public ShardingContexts getJobShardingContext(final List<Integer> shardingItems) {
LiteJobConfiguration liteJobConfig = configService.load(false);
removeRunningIfMonitorExecution(liteJobConfig.isMonitorExecution(), shardingItems);
if (shardingItems.isEmpty()) {
return new ShardingContexts(buildTaskId(liteJobConfig, shardingItems), liteJobConfig.getJobName(), liteJobConfig.getTypeConfig().getCoreConfig().getShardingTotalCount(),
liteJobConfig.getTypeConfig().getCoreConfig().getJobParameter(), Collections.<Integer, String>emptyMap());
}
Map<Integer, String> shardingItemParameterMap = new ShardingItemParameters(liteJobConfig.getTypeConfig().getCoreConfig().getShardingItemParameters()).getMap();
return new ShardingContexts(buildTaskId(liteJobConfig, shardingItems), liteJobConfig.getJobName(), liteJobConfig.getTypeConfig().getCoreConfig().getShardingTotalCount(),
liteJobConfig.getTypeConfig().getCoreConfig().getJobParameter(), getAssignedShardingItemParameterMap(shardingItems, shardingItemParameterMap));
}
private void removeRunningIfMonitorExecution(final boolean monitorExecution, final List<Integer> shardingItems) {
if (!monitorExecution) {
return;
}
List<Integer> runningShardingItems = new ArrayList<>(shardingItems.size());
for (int each : shardingItems) {
if (isRunning(each)) {
runningShardingItems.add(each);
}
}
shardingItems.removeAll(runningShardingItems);
}
正常作业执行的时候会将参数封装为分片上下文ShardingContexts:
taskId 作业任务ID。
jobName 作业名称。
shardingTotalCount 分片总数。
jobParameter 作业自定义参数。
shardingItemParameters 分配于本作业实例的分片项和参数的Map。




