Elastic-Job 实战案例详解
项目简介
Elastic-Job 是当当开源的一套分布式任务调度框架,基于 Quartz 开发,使用 Zookeeper 作为注册中心,实现任务分片、弹性扩容、失效转移等特性。

实战案例:订单超时自动取消
以电商系统中"订单超时自动取消"功能为例,演示 Elastic-Job 的完整使用。
环境准备
<!-- pom.xml 依赖 -->
<dependency>
<groupId>com.dangdang</groupId>
<artifactId>elastic-job-lite-core</artifactId>
<version>2.1.5</version>
</dependency>
<dependency>
<groupId>com.dangdang</groupId>
<artifactId>elastic-job-lite-spring</artifactId>
<version>2.1.5</version>
</dependency>
<!-- Zookeeper -->
<dependency>
<groupId>org.apache.curator</groupId>
<artifactId>curator-framework</artifactId>
<version>4.2.0</version>
</dependency>
配置文件 application.yml
spring:
datasource:
url: jdbc:mysql://localhost:3306/order_db
username: root
password: 123456
driver-class-name: com.mysql.jdbc.Driver
# Elastic-Job 配置
elasticjob:
zookeeper:
server-lists: localhost:2181
namespace: elastic-job
task:
order-cancel:
cron: 0 0/5 * * * ? # 每5分钟执行一次
sharding-total-count: 4 # 总分片数
sharding-item-parameters: 0=beijing,1=shanghai,2=guangzhou,3=shenzhen
failover: true # 故障转移
创建订单实体类
@Data
public class Order {
private Long id;
private String orderNo;
private Integer status; // 0-待支付 1-已支付 2-已取消
private Date createTime;
private Date paymentTime;
private String shardingParam; // 分片参数
// 订单超时时间(分钟)
public static final int TIMEOUT_MINUTES = 30;
}
订单 Mapper 接口
@Mapper
public interface OrderMapper {
// 查询超时订单(根据分片参数)
@Select("SELECT * FROM t_order WHERE status = 0 AND create_time < DATE_SUB(NOW(), INTERVAL 30 MINUTE) AND sharding_param = #{shardingParam} LIMIT 100")
List<Order> selectTimeoutOrders(@Param("shardingParam") String shardingParam);
// 批量更新订单状态为已取消
@Update("<script>" +
"UPDATE t_order SET status = 2, cancel_time = NOW() " +
"WHERE id IN " +
"<foreach collection='ids' item='id' open='(' separator=',' close=')'>" +
"#{id}" +
"</foreach>" +
"</script>")
int batchCancelOrders(@Param("ids") List<Long> ids);
// 更新订单状态
@Update("UPDATE t_order SET status = #{status} WHERE order_no = #{orderNo}")
int updateStatus(@Param("orderNo") String orderNo, @Param("status") Integer status);
}
自定义分片策略(可选)
public class OrderShardingStrategy implements JobShardingStrategy {
@Override
public Map<JobInstance, List<Integer>> sharding(List<JobInstance> jobInstances,
List<Integer> shardingItems) {
Map<JobInstance, List<Integer>> result = new HashMap<>();
// 按 IP 地址做简单轮询分配
for (int i = 0; i < jobInstances.size(); i++) {
JobInstance instance = jobInstances.get(i);
List<Integer> assignedItems = new ArrayList<>();
for (int j = i; j < shardingItems.size(); j += jobInstances.size()) {
assignedItems.add(shardingItems.get(j));
}
result.put(instance, assignedItems);
}
return result;
}
}
任务实现类
@Component
public class OrderCancelJob implements SimpleJob {
private static final Logger log = LoggerFactory.getLogger(OrderCancelJob.class);
@Autowired
private OrderMapper orderMapper;
@Override
public void execute(ShardingContext shardingContext) {
// 获取分片参数
int shardingItem = shardingContext.getShardingItem();
String shardingParam = shardingContext.getShardingParameter();
log.info("=== 订单取消任务开始执行,分片项:{},分片参数:{},JobName:{} ===",
shardingItem, shardingParam, shardingContext.getJobName());
// 查询当前分片的超时订单
List<Order> timeoutOrders = orderMapper.selectTimeoutOrders(shardingParam);
if (timeoutOrders == null || timeoutOrders.isEmpty()) {
log.info("分片{}没有需要处理的超时订单", shardingItem);
return;
}
log.info("分片{}发现{}条超时订单", shardingItem, timeoutOrders.size());
// 批量取消订单
List<Long> orderIds = timeoutOrders.stream()
.map(Order::getId)
.collect(Collectors.toList());
try {
int count = orderMapper.batchCancelOrders(orderIds);
log.info("分片{}成功取消{}条订单", shardingItem, count);
// 记录处理日志
timeoutOrders.forEach(order ->
log.info("订单{}已取消,取消时间:{}",
order.getOrderNo(), new Date())
);
} catch (Exception e) {
log.error("分片{}取消订单失败", shardingItem, e);
throw new RuntimeException("取消订单失败", e);
}
}
}
任务配置类
@Configuration
public class ElasticJobConfig {
@Autowired
private OrderCancelJob orderCancelJob;
@Bean
public CoordinatorRegistryCenter regCenter() {
// 配置 ZooKeeper 注册中心
ZookeeperConfiguration zkConfig = new ZookeeperConfiguration(
new ZookeeperConfiguration.CoordinatorRegistryCenterBuilder()
.serverLists("localhost:2181")
.namespace("elastic-job")
.build()
);
return new ZookeeperRegistryCenter(zkConfig);
}
@Bean
public JobScheduler orderCancelJobScheduler() {
// 配置任务详情
JobCoreConfiguration coreConfig = JobCoreConfiguration.newBuilder(
"orderCancelJob", // 任务名称
"0 0/5 * * * ?", // cron 表达式
4) // 分片总数
.shardingItemParameters("0=beijing,1=shanghai,2=guangzhou,3=shenzhen")
.shardingStrategyClass(OrderShardingStrategy.class.getName())
.failover(true)
.misfire(true)
.description("订单超时自动取消任务")
.build();
SimpleJobConfiguration jobConfig = new SimpleJobConfiguration(
coreConfig,
OrderCancelJob.class.getCanonicalName()
);
LiteJobConfiguration liteJobConfig = LiteJobConfiguration
.newBuilder(jobConfig)
.overwrite(true) // 本地配置覆盖注册中心配置
.build();
return new JobScheduler(regCenter(), liteJobConfig, orderCancelJob);
}
}
数据分片初始化
@Component
public class OrderShardingInitializer implements CommandLineRunner {
@Autowired
private JdbcTemplate jdbcTemplate;
@Override
public void run(String... args) throws Exception {
// 为每个订单分配分片参数
String[] shardingParams = {"beijing", "shanghai", "guangzhou", "shenzhen"};
// 批量更新现有订单的分片字段
for (int i = 0; i < shardingParams.length; i++) {
jdbcTemplate.update(
"UPDATE t_order SET sharding_param = ? WHERE id % 4 = ?",
shardingParams[i], i
);
}
log.info("订单分片初始化完成");
}
}
监控与运维
@RestController
@RequestMapping("/job")
public class JobMonitorController {
@Autowired
private JobOperateAPI jobOperateAPI;
// 查看任务状态
@GetMapping("/status")
public List<JobStatusDTO> getJobStatus() {
return jobOperateAPI.getJobStatus("orderCancelJob");
}
// 手动触发任务
@PostMapping("/trigger")
public void triggerJob() {
jobOperateAPI.trigger("orderCancelJob");
}
// 禁用任务
@PostMapping("/disable")
public void disableJob() {
jobOperateAPI.disable("orderCancelJob", null);
}
// 启用任务
@PostMapping("/enable")
public void enableJob() {
jobOperateAPI.enable("orderCancelJob", null);
}
// 动态修改分片数
@PutMapping("/sharding")
public void updateShardingCount(@RequestParam int count) {
jobOperateAPI.updateShardingCount("orderCancelJob", count);
}
}
任务幂等性处理(防重复执行)
@Component
public class IdempotentUtils {
@Autowired
private StringRedisTemplate redisTemplate;
// 分布式锁 + 幂等性校验
public boolean acquireLock(String jobName, int shardingItem, String jobParameter) {
String key = String.format("job:lock:%s:%s", jobName, shardingItem);
// 尝试获取分布式锁,设置过期时间为5分钟
Boolean locked = redisTemplate.opsForValue()
.setIfAbsent(key, jobParameter, Duration.ofMinutes(5));
if (Boolean.TRUE.equals(locked)) {
log.info("获取任务锁成功:{}", key);
return true;
}
// 检查是否同一参数
String existingParam = redisTemplate.opsForValue().get(key);
if (jobParameter.equals(existingParam)) {
log.warn("任务已在执行中:{}", key);
return false;
}
return false;
}
public void releaseLock(String jobName, int shardingItem) {
String key = String.format("job:lock:%s:%s", jobName, shardingItem);
redisTemplate.delete(key);
}
}
业务服务调用示例
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
// 创建订单
public void createOrder(Order order) {
order.setStatus(0); // 待支付
order.setCreateTime(new Date());
// 设置分片参数
int shardingNum = (int) (order.getId() % 4);
String[] shardingParams = {"beijing", "shanghai", "guangzhou", "shenzhen"};
order.setShardingParam(shardingParams[shardingNum]);
orderMapper.insert(order);
}
// 支付订单
public void payOrder(String orderNo) {
orderMapper.updateStatus(orderNo, 1); // 已支付
}
// 手动取消订单
public void cancelOrder(String orderNo) {
orderMapper.updateStatus(orderNo, 2); // 已取消
}
}
关键特性演示
动态分片调度(流量倾斜处理)
// 动态计算每个分片的压力
@Component
public class DynamicShardingStrategy {
@Autowired
private JdbcTemplate jdbcTemplate;
// 根据数据量动态调整分片策略
public Map<Integer, String> calculateShardingParams() {
Map<Integer, String> result = new HashMap<>();
// 对各区域订单量进行统计
String sql = "SELECT sharding_param, COUNT(*) FROM t_order " +
"WHERE status = 0 GROUP BY sharding_param";
List<Map<String, Object>> stats = jdbcTemplate.queryForList(sql);
// 动态调整分片参数(示例:根据并发量调整)
for (Map<String, Object> stat : stats) {
int count = ((Number) stat.get("COUNT(*)")).intValue();
// 如果订单量很大,可以将参数改为更精细的粒度
if (count > 10000) {
result.put(Integer.parseInt(stat.get("sharding_param").toString()), "high");
} else {
result.put(Integer.parseInt(stat.get("sharding_param").toString()), "normal");
}
}
return result;
}
}
多实例部署测试
# application-instance1.yml(实例1)
elasticjob:
zookeeper:
server-lists: localhost:2181
namespace: elastic-job
instance:
ip: 192.168.1.100
port: 8080
# application-instance2.yml(实例2)
elasticjob:
zookeeper:
server-lists: localhost:2181
namespace: elastic-job
instance:
ip: 192.168.1.101
port: 8081
@SpringBootApplication
public class ElasticJobApplication {
public static void main(String[] args) {
SpringApplication.run(ElasticJobApplication.class, args);
// 启动时自动注册到 ZooKeeper
System.out.println("Elastic-Job 任务调度应用启动成功");
}
// 实例启动时自动注册
@Component
public class JobStartupListener {
@EventListener(ApplicationReadyEvent.class)
public void onStartup() {
JobInstance jobInstance = new JobInstance();
jobInstance.setIp(getLocalIp());
jobInstance.setPort(getPort());
log.info("实例启动:{}", jobInstance);
}
}
}
任务异常处理和重试机制
@Component
public class OrderCancelJobWithRetry implements SimpleJob {
@Autowired
private OrderMapper orderMapper;
@Autowired
private RetryTemplate retryTemplate;
@Override
public void execute(ShardingContext shardingContext) {
int shardingItem = shardingContext.getShardingItem();
String shardingParam = shardingContext.getShardingParameter();
// 使用重试机制处理失败任务
retryTemplate.execute(context -> {
try {
// 业务处理
List<Order> timeoutOrders = orderMapper.selectTimeoutOrders(shardingParam);
if (timeoutOrders != null && !timeoutOrders.isEmpty()) {
List<Long> ids = timeoutOrders.stream()
.map(Order::getId)
.collect(Collectors.toList());
orderMapper.batchCancelOrders(ids);
// 记录处理结果
log.info("分片{}处理完成,取消{}条订单", shardingItem, ids.size());
}
return true;
} catch (Exception e) {
log.error("分片{}处理失败", shardingItem, e);
throw e;
}
}, new RetryCallback<Boolean, Exception>() {
@Override
public Boolean doWithRetry(RetryContext context) throws Exception {
log.warn("重试第{}次", context.getRetryCount() + 1);
return true;
}
});
}
}
最佳实践总结
| 实践项 | 推荐方案 | 说明 |
|---|---|---|
| 分片策略 | 根据数据特点选择 | 订单类任务按取模,用户类按ID范围 |
| 任务幂等 | Redis锁 + 状态检查 | 防止重复执行 |
| 异常处理 | 重试机制 + 告警通知 | 保证任务最终一致性 |
| 性能优化 | 批量操作 + 分页查询 | 避免一次性处理过多数据 |
| 监控体系 | 定时检查 + 日志记录 | 快速定位问题 |
| 配置管理 | 通过配置中心动态调整 | 支持运行时修改 |
常见问题及解决方案
- 任务执行顺序混乱 → 设置任务参数或使用幂等性校验
- 分片不均 → 自定义分片策略,根据数据量动态调整
- 任务失败堆积 → 设置 misfire(true),启用失效转移
- ZooKeeper 宕机 → 使用高可用集群部署
- 任务间依赖 → 使用事件监听器或消息队列
通过以上案例,可以系统性地使用 Elastic-Job 完成分布式任务调度,并解决实际业务场景中的问题。