Elastic-Job 分布式调度
概述
Elastic-Job 是当当网开源的分布式任务调度框架,基于 Zookeeper 实现任务高可用和分片调度。它的核心设计理念是「任务分片 + 弹性扩缩容」。
Elastic-Job 后来捐给了 Apache 基金会,发展为 Apache ShardingSphere ElasticJob。
核心特性
- 分片调度 — 将任务拆分为多个分片,由不同执行器并行处理
- 弹性扩缩容 — 执行器增减时自动重新分片,无需重启
- 高可用 — ZK 选举主节点,主节点负责分片分配
- 运维平台 — Web 控制台可视化任务管理
一、架构
┌─────────────────────────────┐
│ ZooKeeper │
│ ┌─────────────────────┐ │
│ │ /elastic-job │ │
│ │ ├── jobA │ │
│ │ │ ├── instances │ │
│ │ │ ├── sharding │ │
│ │ │ └── leader │ │
│ │ └── jobB │ │
│ └─────────────────────┘ │
└──────────┬──────────────────┘
│
┌──────┼──────┐
▼ ▼ ▼
┌──────┐┌──────┐┌──────┐
│执行器1││执行器2││执行器3│
│分片0 ││分片1 ││分片2 │
│分片3 ││ ││ │
└──────┘└──────┘└──────┘二、Spring Boot 集成
2.1 依赖
xml
<dependency>
<groupId>org.apache.shardingsphere.elasticjob</groupId>
<artifactId>elasticjob-lite-spring-boot-starter</artifactId>
<version>3.0.3</version>
</dependency>2.2 配置
yaml
elasticjob:
reg-center:
server-lists: localhost:2181
namespace: elastic-job
jobs:
stockSyncJob:
elastic-job-class: com.example.job.StockSyncJob
cron: 0 0/5 * * * ?
sharding-total-count: 3
sharding-item-parameters: 0=A,1=B,2=C
props:
print: true2.3 任务实现
java
@Component
public class StockSyncJob implements SimpleJob {
@Autowired
private StockService stockService;
@Override
public void execute(ShardingContext context) {
int shardItem = context.getShardingItem(); // 0, 1, 2
String shardParam = context.getShardingParameter(); // A, B, C
log.info("分片 {} (参数 {}) 开始同步库存", shardItem, shardParam);
stockService.syncStock(shardParam);
}
}三、分片策略
3.1 内置分片策略
| 策略类 | 说明 |
|---|---|
AverageAllocationJobShardingStrategy | 平均分配(默认) |
OdevitySortByNameJobShardingStrategy | 奇偶分组 |
RotateServerByNameJobShardingStrategy | 轮询分配 |
3.2 自定义分片
java
public class WarehouseShardingStrategy implements JobShardingStrategy {
@Override
public Map<Integer, String> sharding(
List<JobInstance> jobInstances, String jobName, int shardingTotalCount) {
Map<Integer, String> result = new LinkedHashMap<>();
// 按权重分配:执行器性能不同时使用
for (int i = 0; i < shardingTotalCount; i++) {
JobInstance instance = jobInstances.get(i % jobInstances.size());
result.put(i, instance.getJobParameter());
}
return result;
}
}四、任务监听
4.1 任务生命周期监听
java
@Component
public class StockSyncListener implements ElasticJobListener {
@Override
public void beforeJobExecuted(ShardingContexts shardingContexts) {
log.info("任务 {} 开始执行,分片数: {}",
shardingContexts.getJobName(),
shardingContexts.getShardingTotalCount());
}
@Override
public void afterJobExecuted(ShardingContexts shardingContexts) {
log.info("任务 {} 执行完成", shardingContexts.getJobName());
}
}4.2 分布式监听(主节点回调)
java
public class DistributedJobListener extends AbstractDistributeOnceElasticJobListener {
public DistributedJobListener(long startedTimeout, long completedTimeout) {
super(startedTimeout, completedTimeout);
}
@Override
public void doBeforeJobExecutedAtLastStarted(ShardingContexts shardingContexts) {
// 仅在最后一个分片开始前执行(适合预检查)
log.info("所有分片即将开始,执行前置检查");
}
@Override
public void doAfterJobExecutedAtLastCompleted(ShardingContexts shardingContexts) {
// 仅在最后一个分片完成后执行(适合后置汇总)
log.info("所有分片执行完成,执行后置处理");
}
}五、事件追踪与运维
5.1 事件追踪
java
@Component
public class JobEventConfig {
@Bean
public JobEventConfiguration jobEventConfiguration() {
// 保存到数据库
return new JobEventRdbConfiguration(dataSource);
}
}5.2 运维界面
Elastic-Job 提供 Web 控制台:
yaml
# elastic-job-lite-console 部署
server:
port: 8088
elasticjob:
reg-center:
server-lists: localhost:2181
namespace: elastic-job控制台功能:
- 任务启停
- 分片监控
- 执行历史
- 事件追踪查看
六、实战:库存同步 + 动态扩缩容
6.1 场景
电商系统需要将库存数据从 Redis 同步到 MySQL,库存数据按仓库编码(A/B/C/D)分片,服务实例支持动态扩缩容。
6.2 任务实现
java
@Component
@ElasticJobConfig(
name = "stockSyncJob",
cron = "0 0/1 * * * ?",
shardingTotalCount = 4,
shardingItemParameters = "0=A,1=B,2=C,3=D",
jobStrategy = AverageAllocationJobShardingStrategy.class
)
public class StockSyncJob implements SimpleJob {
@Override
public void execute(ShardingContext context) {
String warehouseCode = context.getShardingParameter();
log.info("仓库 {} 开始同步库存", warehouseCode);
// 从 Redis 读取该仓库的库存
Map<String, Integer> stockMap = redisTemplate
.opsForHash().entries("stock:" + warehouseCode);
// 批量写入 MySQL
for (Map.Entry<String, Integer> entry : stockMap.entrySet()) {
stockDao.updateStock(warehouseCode, entry.getKey(), entry.getValue());
}
log.info("仓库 {} 同步完成,共 {} 条", warehouseCode, stockMap.size());
}
}6.3 动态扩缩容测试
初始:3 个实例,4 个分片
┌────────┐ ┌────────┐ ┌────────┐
│实例1 │ │实例2 │ │实例3 │
│分片0(A) │ │分片1(B) │ │分片2(C) │
│分片3(D) │ │ │ │ │
└────────┘ └────────┘ └────────┘
扩容:新增实例4 → 自动重新分配
┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐
│实例1 │ │实例2 │ │实例3 │ │实例4 │
│分片0(A) │ │分片1(B) │ │分片2(C) │ │分片3(D) │
└────────┘ └────────┘ └────────┘ └────────┘
缩容:下线实例2 → 自动重新分配
┌────────┐ ┌────────┐ ┌────────┐
│实例1 │ │实例3 │ │实例4 │
│分片0(A) │ │分片1(B) │ │分片2(C) │
│分片3(D) │ │ │ │ │
└────────┘ └────────┘ └────────┘七、@Scheduled / XXL-Job / Elastic-Job 对比
| 对比维度 | @Scheduled | XXL-Job | Elastic-Job |
|---|---|---|---|
| 架构 | 单体内置 | 调度中心 + 执行器 | ZK + 执行器 |
| 分片 | 不支持 | 分片广播 | 原生分片调度 |
| 高可用 | 无(多实例重复执行) | 失败转移 | ZK 选举 + 自动重分片 |
| 动态扩缩容 | 不支持 | 手动调整 | 自动感知 re-sharding |
| 可视化控制台 | 无 | 完善 | 完善 |
| 学习成本 | 低 | 中 | 中高 |
| 适合场景 | 简单单机定时 | 企业级分布式任务 | 大规模数据分片处理 |
选型决策树
是否需要集群/高可用?
├── 否 → @Scheduled
└── 是 → 是否需要分片?
├── 是 → 是否需要弹性扩缩容?
│ ├── 是 → Elastic-Job
│ └── 否 → XXL-Job
└── 否 → XXL-Job八、总结
| 知识点 | 说明 |
|---|---|
| 架构 | ZK + 执行器,主节点负责分片分配 |
| 分片策略 | 平均分配 / 奇偶分组 / 轮询 / 自定义 |
| 弹性扩缩容 | 实例增减自动重分片,无需重启 |
| 监听 | ElasticJobListener 生命周期 + DistributeOnce 主节点回调 |
| 事件追踪 | 支持 RDB / 日志追踪 |
| 与 XXL-Job 对比 | ES-Job 分片能力更强,XXL-Job 部署更简单 |
参考链接: