Elastic-Job 与 PowerJob 源码阅读
Elastic-Job 以分片见长,靠 Zookeeper 自协调;PowerJob 以工作流 DAG 与秒级调度见长,采用 Server/Worker 中心化架构。本文分别拆解两者的核心源码。
Elastic-Job 源码
整体架构
Elastic-Job 节点构成:
├─ Job:任务(SimpleJob / DataflowJob / ScriptJob)
├─ RegistryCenter:注册中心(Zookeeper),存放作业与实例元数据
└─ JobScheduler:作业调度器(整合 Quartz 做时间触发)
无独立调度中心:每个作业节点都是"调度者 + 执行者"Zookeeper 数据节点
ZK 中作业的数据结构(/jobName/...):
/jobName/instances 实例节点(各执行节点的临时节点)
/jobName/sharding 分片节点({分片号}/instance 指定谁执行)
/jobName/leader 选举节点(/election/instance)
/jobName/config 作业配置(cron、分片数等)
/jobName/servers 注册的服务器节点
协作方式:
├─ 实例上线 → 创建临时节点 → 触发重新分片
├─ 实例下线 → 临时节点消失 → 触发重新分片
└─ 通过 ZK 的 Watch 机制感知变化分片流程源码
分片触发流程:
1. 作业执行前,检查是否需要重新分片(shardingNode)
2. 触发分片 → 由"分片主节点"执行分片计算
3. 分片结果写入 ZK 的 /sharding 节点
4. 各实例读取自己应执行的分片分片主节点选举
java
// 作业启动时通过 leader 选举选出一个"分片主节点"
// LeaderElectionService
public class LeaderElectionService {
public void electLeader() {
// 创建临时顺序节点参与选举
jobNodeStorage.createJobNodeIfNeeded(LeaderNode.ELECTION_LATCH);
// 获取最小序号节点 → 胜出者成为 leader
List<String> nodes = jobNodeStorage.getJobNodeChildrenKeys(ELECTION_LATCH);
Collections.sort(nodes);
if (candidate.equals(nodes.get(0))) {
// 自己是 leader → 写入 leader 节点
jobNodeStorage.fillEphemeralJobNode(LeaderNode.INSTANCE, ...);
}
}
}分片计算(ShardingService)
java
public class ShardingService {
// 分片主节点执行
public void shardingIfNecessary() {
if (jobNodeStorage.isJobNodeExisted(ShardingNode.NECESSARY)) {
// 1. 抢占"分片处理中"标志(防止并发分片)
if (jobNodeStorage.fillEphemeralJobNode(SHARDING_FLAG, ...)) {
// 2. 收集所有可用实例
List<JobInstance> availableJobInstances = instanceService.getAvailableJobInstances();
// 3. 按分片策略分配
Map<JobInstance, List<Integer>> result = shardingStrategy.sharding(
availableJobInstances, jobConfig.getShardingTotalCount());
// 4. 把分配结果写入 ZK 分片节点
// 5. 清除"需要分片"标志
}
}
}
}分片策略
分片策略接口 ShardingStrategy:
├─ AverageAllocationShardingStrategy:平均分配(默认)
│ 分片 0..n 轮流分配给实例
├─ OdevitySortByNameJobShardingStrategy:按名称哈希分配
├─ RotateServerByNameJobShardingStrategy:轮询
└─ 自定义策略:实现接口注册
平均分配示例(5 片,2 实例):
├─ 实例A:分片 0、1、2、3
└─ 实例B:分片 4作业执行与分片上下文
java
// 作业接口:SimpleJob
public interface SimpleJob {
void execute(ShardingContext shardingContext);
}
// 分片上下文:执行时拿到自己的分片
public final class ShardingContext {
private int shardingItem; // 当前分片号
private int shardingTotalCount; // 总分片数
private String jobName;
private String jobParameter; // 任务参数
}作业执行流程
作业调度执行:
1. JobScheduler 内嵌 Quartz,按 cron 触发
2. 触发后执行分片检查(是否需要重新分片)
3. 从 ZK 读取当前实例的分片列表
4. 每个分片调用一次 execute(ShardingContext)
5. 执行期间通过监听器上报状态高可用与失效转移
失效转移(failover):
├─ 实例执行中挂掉 → 其分片变成"未分配"
├─ 监听器检测到实例下线 → 触发重新分片
├─ 存活的实例接管失效分片
└─ 配合 Misfire(错过触发补偿):错过的调度补跑
幂等保障:
├─ 作业可配 misfire 开关(错过是否补偿)
└─ 业务侧结合分片号做幂等处理DataflowJob 流式处理
java
// 流式作业:抓取一批 → 处理一批 → 循环
public interface DataflowJob<T> {
List<T> fetchData(ShardingContext shardingContext); // 抓取
void processData(ShardingContext shardingContext, List<T> data); // 处理
}流式执行:
├─ 不断 fetchData → processData
├─ fetchData 返回空 → 本轮结束(等下次调度)
└─ 适合:数据搬运、分批处理大批量数据PowerJob 源码
整体架构
PowerJob 四组件:
├─ Server:调度中心(集群部署,共享 DB)
│ ├─ 调度引擎(cron / 秒级 / 工作流触发)
│ ├─ Worker 管理(心跳、路由)
│ └─ 任务状态维护
├─ Worker:执行器(内嵌业务应用)
│ ├─ 执行任务处理器
│ └─ 回传执行状态与日志
├─ Akka:Server 与 Worker 的通信通道
└─ 存储:MySQL(元数据)+ MongoDB(日志)Server 调度核心
调度流程:
1. 调度线程扫描到期任务
2. 按任务配置选 Worker(路由策略)
3. 通过 Akka 发送任务执行请求
4. Worker 执行并回传状态
5. Server 更新任务实例状态工作流 DAG
工作流 = 有向无环图(DAG):
├─ 节点:任务(一个任务可被多个任务依赖)
├─ 边:依赖关系(前驱完成后触发后继)
└─ 支持并行分支与汇合
示例:订单数据处理工作流
├─ 节点A:拉取订单(先)
├─ 节点B:清洗数据(依赖 A)
├─ 节点C:计算指标(依赖 A)
└─ 节点D:生成报表(依赖 B、C,即汇合)DAG 执行引擎源码思路
java
// 工作流执行引擎:拓扑排序驱动
public class WorkflowDAG {
// 节点执行状态
public void tryTriggerNodes(...) {
for (WorkflowNode node : dag.getNodes()) {
// 1. 检查该节点是否满足执行条件
// (所有前驱节点都已完成)
boolean ready = node.getDependencies().stream()
.allMatch(dep -> isFinished(dep));
if (ready && !isStarted(node)) {
// 2. 满足条件 → 发起执行
dispatch(node);
}
}
// 3. 全部节点完成 → 工作流结束
}
}核心机制:
├─ 依赖完成事件驱动(前驱完成 → 通知后继检查)
├─ 并行分支:多个无依赖关系的节点同时执行
├─ 汇合:后继等待所有前驱完成
└─ 失败处理:节点失败 → 可配置中断整个工作流或跳过后继秒级调度
PowerJob 支持秒级任务:
├─ 不同于 cron 的分钟精度
├─ 支持 1s / 5s / 10s 固定间隔触发
└─ 实现:Server 维护时间轮,高精度触发MapReduce 分布式计算
MapReduce 任务模型:
├─ 任务拆分成多个子任务(Map)
├─ 子任务分发到多个 Worker 并行执行
├─ 执行结果聚合(Reduce)
└─ 适合:大数据量分布式处理(类似数据批处理)
示例:全量订单统计
├─ Map:按订单号段拆成 N 个子任务
├─ 每个 Worker 计算自己子任务的汇总
└─ Reduce:合并各 Worker 的汇总 → 最终结果两者对比
| 对比项 | Elastic-Job | PowerJob |
|---|---|---|
| 架构 | ZK 自协调(无中心) | Server + Worker(中心化) |
| 调度精度 | 分钟级 | 秒级 |
| 分片 | 强(核心特性) | 有 |
| 工作流 DAG | 无 | 有 |
| 分布式计算 | 无 | MapReduce |
| 依赖组件 | Zookeeper | MySQL + MongoDB |
| 运维界面 | 需自建 | 自带控制台 |
| 定位 | 分片并行作业 | 分布式计算与编排 |
选型提示:
├─ 只要分片并行、已有 ZK → Elastic-Job
├─ 需要秒级调度、DAG 编排 → PowerJob
└─ 常规任务管理 → XXL-Job 更轻总结
Elastic-Job 的源码核心是 ZK 数据节点 + 分片主节点 + 分片策略:所有实例通过 ZK 感知彼此,由 leader 统一计算分片并写入 ZK,各实例各司其职。PowerJob 的源码核心是 Server 调度引擎 + DAG 工作流引擎 + MapReduce:Server 统一调度与状态管理,工作流按依赖关系驱动节点执行,秒级调度与分布式计算是它的差异化能力。