Elastic-Job案例

wen java案例 1

Elastic-Job 实战案例详解

项目简介

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

Elastic-Job案例

实战案例:订单超时自动取消

以电商系统中"订单超时自动取消"功能为例,演示 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锁 + 状态检查 防止重复执行
异常处理 重试机制 + 告警通知 保证任务最终一致性
性能优化 批量操作 + 分页查询 避免一次性处理过多数据
监控体系 定时检查 + 日志记录 快速定位问题
配置管理 通过配置中心动态调整 支持运行时修改

常见问题及解决方案

  1. 任务执行顺序混乱 → 设置任务参数或使用幂等性校验
  2. 分片不均 → 自定义分片策略,根据数据量动态调整
  3. 任务失败堆积 → 设置 misfire(true),启用失效转移
  4. ZooKeeper 宕机 → 使用高可用集群部署
  5. 任务间依赖 → 使用事件监听器或消息队列

通过以上案例,可以系统性地使用 Elastic-Job 完成分布式任务调度,并解决实际业务场景中的问题。

抱歉,评论功能暂时关闭!