一、解决什么问题
一个任务中心技术实现的参考案例,可以快速部署实现且仅需关注业务个性落库逻辑实现,其他如任务状态维护、数据解析及异常包装、结果导出均由工具自动实现。
二、基本原理

图1 请求示意图
异步任务中心共分三个模块:
1)任务初始化, 将目标导入文件上传至云存储后得到目标文件url按任务类型(如类目导入、商品导入等)入库任务表并返回前台提交成功,任务初始状态为"待处理";
2)任务调度,使用开源调度组件xxlJob开箱即用。传送门: xxlJob
3)任务Worker执行器核心组成:
1.任务并行分片拉取
分片广播模式下,每个worker按index取模 获取应执行的任务id,参考sql :
from task where status in ('PENDING','FAILURE') and errCnt <= MAX_RETRY_CNT and mod(id,#总worker数量) = 当前worker index
2.根据任务类型命中执行器策略
任务类型: 即导入业务的枚举字段,如类目导入CATE_IMPORT、商品导入PRODUCT_IMPORT等
业务执行器: 执行excel批量导入解析落库的载体,下文介绍。
策略如何命中: 业务执行器class类增加@JobExecutor注解并指明注解值为对应任务类型; 拉取任务后寻找有@JobExecutor修饰的类且其注解值等于任务记录任务类型即为命中目标执行器
3.执行器设计
A.抽象任务接口并定义行为 -> BaseJob<T>
•accept() 接受任务,实现后置任务状态为"处理中"
•parse() 解析任务, 负责解析目标文件(zip、xlsx)为List<Bean>,并实现数据校验
•run() 将业务数据List<Bean>数据落库
•export() 生成导入结果文件,上传至云存储并更新到任务记录结果列
•errHandle() 异常处理,置任务状态为"失败",累计任务失败次数,触发业务报警
B.基础抽象实现类 -> BaseExecutableAbsJob implements BaseJob
accept() 、export()、errHandle() 步骤因其业务无关性故在此抽象类中做通用默认实现;
parse() 有一定通用性,默认实现为excel解析(easyExcel实现)
run() 业务相关不做默认实现,由继承方实现
C.一次性解析抽象实现 -> DisposableAbsJob extends BaseExecutableAbsJob
特征:
解析规则为一次性解析excel所有记录**,不适用超大excel解析job**
可以在落库前获得全部业务实体信息
导出结果可以显示原始输入
D.分批解析通用实现类 -> BatchableAbsJob extends BaseExecutableAbsJob
特征:
解析规则为按BATCH_CNT来分段操作数据解析及入库,适用于大excel导入场景的使用
解析完毕前拿不到记录总数
导出结果不显示原始输入,仅显示MAX_ERROR_CNT数量以内的错误记录原始信息及错误信息。
三、快速使用
业务类按场景选择继承DisposableAbsJob或BatchableAbsJob,
**仅需重写落库方法,其他如拉取、解析、导出结果步骤均由系统自动执行。如需特殊解析逻辑(比如解析zip按特定规则拼装bean)重写parse()**方法即可
举个栗子,现需求场景为批量类目信息导入, 则开发过程为:
步骤一 : 落库任务类型为TaskBizTypeEnum.CATE_BATCH_PUBLISH的记录到任务表中,并记录前台上传的excel导入文件url(常规crud本案例不做封装,自行实现即可)
步骤二 : 定义类目Excel导入实体Bean
1/** 2 * 类目导入实体 3 */ 4@Data 5@NoArgsConstructor 6@AllArgsConstructor 7@EqualsAndHashCode(callSuper = true) 8public class ImportCateExcelDTO extends BaseWorkerDTO { 9 10 /** 类目级别*/ 11 @ExcelProperty(index = 0,converter = CateLevelConverter.class,value = "类目级别") 12 private Integer cateLevel; 13 14 /** 类目中文名*/ 15 @ExcelProperty(index = 1 ,value = "类目中文名") 16 private String cateName; 17 18 /** 类目排序*/ 19 @ExcelProperty(index = 2 ,value = "类目排序") 20 private Integer sort; 21 22 /** 上级类目id*/ 23 @ExcelProperty(index = 3 ,value = "上级管理类目id") 24 private Long parentCateId; 25 26 /** 状态*/ 27 @ExcelProperty(index = 4,converter = StatusConverter.class ,value = "状态") 28 private Integer status; 29 30}
步骤三 : 编写业务实现类,并自行实现run落库方法.
1/** 2 * 类目批量导入(一次性解析全部excel) 3 */ 4@Service 5@Slf4j 6@JobExecutor(taskBizType = TaskBizTypeEnum.CATE_BATCH_PUBLISH) // 策略注解,枚举类型全局唯一。 不加该注解则任务调度找不到策略 7public class DisposableCateImportHandler extends DisposableAbsJob<ImportCateExcelDTO> { 8 9 @Resource 10 private XXXXService xxxxService; 11 12 @Override 13 public void run(TaskDTO<ImportCategoryExcelDTO> task){ 14 try{ 15 if(CollectionUtils.isNotEmpty(task.getTarget())){ 16 xxxxService.save(task.getTarget()) 17 } 18 }catch (BaseImportException e){ 19 errHandle(task); 20 } 21 } 22}
至此开发部分结束,任务执行器会自动调度拉取CATE_BATCH_PUBLISH类型的任务 -> 解析到List<Bean> -> 调用你的run()方法实现落库 -> 将结果流上传到云存储并将结果链接更新到任务表中
四、源码
1. TaskDispatcher - 任务调度派发
1/** 2 * 任务调度派发 3 */ 4@Component 5@Slf4j 6public class TaskDispatcher { 7 8 @Resource 9 private TaskMangeService taskMangeService; 10 @Resource 11 private ApplicationContext applicationContext; 12 13 @SneakyThrows 14 @XxlJob("iscWorker") 15 public ReturnT<String> iscWorker(String param) { 16 TaskDTO task = taskMangeService.pullTask(); 17 if(task!=null){ 18 BaseJob executor = getExecutor(task.getTask().getBizType()); 19 if(null!=executor){ 20 executor.of(task).start(); 21 log.info("iscWorker 执行完毕:{} " , JSON.toJSONString(task)); 22 } 23 } 24 return ReturnT.SUCCESS; 25 } 26 27 //获取执行器 28 public BaseJob getExecutor(TaskBizTypeEnum taskBizType){ 29 Map<String, Object> beanMap = applicationContext.getBeansWithAnnotation(JobExecutor.class); 30 if(beanMap.isEmpty()){ 31 return null; 32 } 33 log.info("TaskDispatcher.getExecutor class list:{}" , beanMap.keySet()); 34 for (Map.Entry<String,Object> entry : beanMap.entrySet()) { 35 try { 36 JobExecutor ano = AnnotationUtil.getAnnotation(entry.getValue().getClass(), JobExecutor.class); 37 if(taskBizType.equals(ano.taskBizType()) && entry.getValue() instanceof BaseJob){ 38 log.info("TaskDispatcher.getExecutor 当前任务:{}命中执行策略job:{}" , taskBizType, entry.getValue()); 39 return (BaseJob) entry.getValue(); 40 } 41 }catch (Exception e){ 42 e.printStackTrace(); 43 } 44 } 45 return null; 46 } 47 48} 49 50
2. DisposableAbsJob - 一次性解析任务执行器
1/** 2 * 一次性解析任务执行器,解析规则为一次性解析所有excel记录,不适用超大excel解析job 3 * 使用方法: 1.使用方继承DisposableAbsJob类,并根据需要重写parse方法(当前默认是按excel解析) 4 * 2.重写run方法,将解析好的list<Bean>推入数据库 5 */ 6@Component 7@Slf4j 8public abstract class DisposableAbsJob<T extends BaseWorkerDTO> extends BaseExecutableAbsJob<T> { 9 10 //自有个性逻辑,默认就是空逻辑 11}
3. BatchableAbsJob - 分段解析任务执行器
1/** 2 * 批次解析任务执行器,解析规则为分批解析excel记录,适用超大excel解析job 3 * 使用方法: 1.使用方继承BatchableAbsJob类,重写saveOrUpdate方法和excel2Po方法, 4 */ 5@Component 6@Slf4j 7public abstract class BatchableAbsJob<T extends BaseWorkerDTO,K> extends BaseExecutableAbsJob<T> { 8 9 /** 10 * 批次解析逻辑 11 * @param task 12 */ 13 @Override 14 public void parse(TaskDTO<T> task){ 15 if(TaskCreateTypeEnum.IMPORT.equals(task.getTaskCreateType())){ 16 log.info("BaseExecutableAbsJob.import parse {} ",task.getTaskId()); 17 BaseBatchExcelDataListener<T,K> listener = new BaseBatchExcelDataListener<>(this); 18 EasyExcel.read(task.getTargetInputFile().getObjectContent(), getTargetClass(), listener).sheet().doRead(); 19 task.setErrDataList(listener.errDataList); 20 } 21 } 22 23 /** 批次解析结果逻辑,仅导出有问题的记录(上限100条) */ 24 @Override 25 public void export(TaskDTO<T> task){ 26 if(task!=null){ 27 log.info("BatchableAbsJob.export {}", task.getTaskId()); 28 if(CollectionUtils.isEmpty(task.getErrDataList())){ 29 taskMangeService.update(new TaskVO(task.getTaskId(), TaskStatusEnum.SUCCESS)); 30 log.info("BatchableAbsJob.export 任务{}全部执行成功" , task.getTaskId()); 31 return; 32 } 33 String resultName = task.getFileName() + Constant.UNDER_LINE + System.currentTimeMillis() + ".xlsx"; 34 ByteArrayOutputStream targetOutputStream = new ByteArrayOutputStream(); 35 try (ExcelWriter excelWriter = EasyExcel.write(targetOutputStream).build()) { 36 if (CollectionUtils.isNotEmpty(task.getErrDataList())) { 37 excelWriter.write(task.getErrDataList(), EasyExcel.writerSheet(0, "result").head(BatchResultDTO.class).build()); 38 } 39 task.setEndTime(System.currentTimeMillis()); 40 excelWriter.finish(); 41 try (ByteArrayInputStream inputStream = new ByteArrayInputStream(targetOutputStream.toByteArray())) { 42 task.setResultUrl(s3Utils.upload(inputStream, FileTypeEnum.BATCH_FILE.getCode(),resultName)); 43 taskMangeService.update(new TaskVO(task.getTaskId(), TaskStatusEnum.SUCCESS, task.getResultUrl())); 44 } 45 } catch (Exception e) { 46 log.error("BaseExecutableAbsJob.export error, target:{} ", task.getTaskId(), e); 47 throw new TaskExportException(task.getTaskId() + e.getMessage()); 48 } finally { 49 log.info("BaseExecutableAbsJob.export 任务「{}」执行完毕:{},文件地址:{}", task.getTaskId(), task.getOssPutMd5(), task.getResultUrl()); 50 } 51 } 52 } 53 54 public List<BatchResultDTO> saveOrUpdate(Map<Integer, K> k) { 55 return null; 56 } 57 58 public Map<Integer,K> excel2Po(Map<Integer, T> excel) { 59 return null; 60 } 61 62}
4. BaseExecutableAbsJob - 通用抽象任务执行器
1/** 2 * 通用抽象任务执行器 3 */ 4@Component 5@Slf4j 6public abstract class BaseExecutableAbsJob<T extends BaseWorkerDTO> implements BaseJob<T> { 7 8 @Resource 9 public S3Utils s3Utils; 10 11 @Resource 12 public TaskMangeService taskMangeService; 13 14 public final static String RESULT_FOLDER = "xxx"; 15 16 17 @Override 18 public void accept(TaskDTO<T> task){ 19 //导入类任务 20 if(TaskCreateTypeEnum.IMPORT.equals(task.getTask().getCreateType())){ 21 task.setTargetInputFile(s3Utils.download(task.getTask().getReqParam())); 22 task.setFileName(task.getTask().getName()); 23 //导出类任务 24 }else if(TaskCreateTypeEnum.EXPORT.equals(task.getTask().getCreateType())){ 25 // 方式1. 保存 前台勾选的记录id到任务入参中 26 // 方式2. 根据前台勾选的查询条件命中记录id,再保存到任务入参中<限制总导出记录数> 27 String req = task.getTask().getReqParam(); 28 if(StringUtils.isNotBlank(req)){ 29 task.setKey(Arrays.stream(req.split(Constant.COMMA)).map(Long::valueOf).collect(Collectors.toSet())); 30 } 31 } 32 task.setTaskBizTypeEnum(task.getTask().getBizType()); 33 task.setTaskId(task.getTask().getId()); 34 task.setStartTime(System.currentTimeMillis()); 35 //更新任务状态 36 taskMangeService.update(new TaskVO(task.getTaskId(),TaskStatusEnum.PROCESSING)); 37 } 38 39 /** 40 * 通用解析逻辑 41 * @param task 42 */ 43 @Override 44 public void parse(TaskDTO<T> task){ 45 if(TaskCreateTypeEnum.IMPORT.equals(task.getTaskCreateType())){ 46 if(task.getTargetInputFile()!=null && task.getTargetInputFile().getObjectContent()!=null){ 47 List<T> target = EasyExcel.read(task.getTargetInputFile().getObjectContent(), getTargetClass() , 48 new PageReadListener<T>(dataList -> {})).sheet(0).headRowNumber(1).doReadSync(); 49 task.setTarget(target); 50 } 51 } 52 } 53 54 /** 55 * 导入通用落库逻辑/导出构建list<Bean>逻辑 56 * @param task 57 */ 58 @Override 59 public void run(TaskDTO<T> task){ } 60 61 @Override 62 public void export(TaskDTO<T> task){ 63 if(task!=null){ 64 if(CollectionUtils.isEmpty(task.getTarget())){ 65 taskMangeService.update(new TaskVO(task.getTaskId(), TaskStatusEnum.SUCCESS)); 66 log.info("BaseExecutableAbsJob.export 空任务{},跳过执行" , task.getTaskId()); 67 return; 68 } 69 String resultName = RESULT_FOLDER + task.getTaskBizTypeEnum().getName() + Constant.UNDER_LINE + System.currentTimeMillis() + ".xlsx"; 70 ByteArrayOutputStream targetOutputStream = new ByteArrayOutputStream(); 71 try (ExcelWriter excelWriter = EasyExcel.write(targetOutputStream).build()) { 72 if (CollectionUtils.isNotEmpty(task.getTarget())) { 73 excelWriter.write(task.getTarget(), EasyExcel.writerSheet(0, "result").head(getTargetClass()).build()); 74 } 75 task.setEndTime(System.currentTimeMillis()); 76 excelWriter.finish(); 77 try (ByteArrayInputStream inputStream = new ByteArrayInputStream(targetOutputStream.toByteArray())) { 78 task.setResultUrl(s3Utils.upload(inputStream, FileTypeEnum.BATCH_FILE.getCode(),resultName)); 79 taskMangeService.update(new TaskVO(task.getTaskId(), TaskStatusEnum.SUCCESS, task.getResultUrl())); 80 } 81 } catch (Exception e) { 82 log.error("BaseExecutableAbsJob.export error, target:{} ", task.getTaskId(), e); 83 throw new TaskExportException(task.getTaskId() + e.getMessage()); 84 } finally { 85 log.info("BaseExecutableAbsJob.export 任务「{}」执行完毕:{},文件地址:{}", task.getTaskId(), task.getOssPutMd5(), task.getResultUrl()); 86 } 87 } 88 } 89 90 @Override 91 public void errHandle(TaskDTO<T> taskDTO,Exception e){ 92 taskMangeService.errHandle(taskDTO,e.toString()); 93 } 94 95 public Class<T> getTargetClass(){ 96 Type res = getClass().getGenericSuperclass(); 97 if(res instanceof ParameterizedType){ 98 ParameterizedType pRes = (ParameterizedType) res; 99 Type[] type = pRes.getActualTypeArguments(); 100 if(type.length>0){ 101 if(type[0] instanceof Class){ 102 Type typeE = type[0]; 103 return (Class<T>)typeE; 104 } 105 } 106 } 107 return null; 108 } 109 110}
5. BaseBatchExcelDataListener - 批处理excel解析监听器
1/** 2 * 批处理excel解析监听器 3 * @param <T> Excel DTO 4 * @param <K> 落库 PO 5 */ 6@Slf4j 7public class BaseBatchExcelDataListener<T extends BaseWorkerDTO,K> implements ReadListener<T> { 8 9 private static final int BATCH_COUNT = 100; 10 private static final int MAX_ERROR_COUNT = 100; 11 12 /** 业务服务*/ 13 private final BatchableAbsJob<T,K> batchableAbsJob; 14 15 /** 每批待处理业务数据*/ 16 private Map<Integer,T> cachedDataList = Maps.newHashMapWithExpectedSize(BATCH_COUNT); 17 18 /** 业务处理失败数据,行号&错误报文 */ 19 public List<BatchResultDTO> errDataList = Lists.newArrayListWithExpectedSize(MAX_ERROR_COUNT) ; 20 21 public BaseBatchExcelDataListener(BatchableAbsJob<T,K> batchableAbsJob) { 22 this.batchableAbsJob = batchableAbsJob; 23 } 24 25 @Override 26 public void invoke(T data, AnalysisContext context) { 27 cachedDataList.put(context.readRowHolder().getRowIndex(),data); 28 if (cachedDataList.size() >= BATCH_COUNT) { 29 saveData(); 30 cachedDataList = Maps.newHashMapWithExpectedSize(BATCH_COUNT); 31 } 32 } 33 34 @Override 35 public void doAfterAllAnalysed(AnalysisContext analysisContext) { 36 saveData(); 37 } 38 39 /** 持久化 */ 40 private void saveData() { 41 Map<Integer, K> po = batchableAbsJob.excel2Po(cachedDataList); 42 if(po!=null && !po.isEmpty()){ 43 List<BatchResultDTO> errRes = batchableAbsJob.saveOrUpdate(po); 44 if(errDataList.size()<MAX_ERROR_COUNT && CollectionUtils.isNotEmpty(errRes)){ 45 errDataList.addAll(errRes); 46 } 47 } 48 } 49} 50 51
6. BaseJob - 任务接口
1public interface BaseJob<T> { 2 3 void accept(TaskDTO<T> task); 4 5 void parse(TaskDTO<T> task); 6 7 void run(TaskDTO<T> task); 8 9 void export(TaskDTO<T> task); 10 11 void errHandle(TaskDTO<T> task,Exception e); 12 13 default AbsExecutor<Void> of(TaskDTO<T> task){ 14 return () -> { 15 try { 16 accept(task); 17 try { 18 parse(task); 19 }finally { 20 if(task.getTargetInputFile()!=null){ 21 task.getTargetInputFile().close(); 22 } 23 } 24 run(task); 25 export(task); 26 }catch (Exception e){ 27 errHandle(task,e); 28 } 29 return null; 30 }; 31 } 32} 33 34
7. JobExecutor- 策略注解
1/** 2 * 任务执行器 3 */ 4@Target({ElementType.TYPE}) 5@Retention(RetentionPolicy.RUNTIME) 6@Documented 7@Inherited 8public @interface JobExecutor { 9 //任务业务类型 10 TaskBizTypeEnum taskBizType() ; 11}
8. TaskMangeService- 任务执行类
1/** 2 * 任务读写服务 3 */ 4@Service 5@Slf4j 6public class TaskMangeServiceImpl extends BaseManageSupportService<TaskVO, TaskPO> implements TaskMangeService { 7 8 private final static Integer MAX_ERR_CNT = 2; 9 private final static Long LIMIT = 1L; 10 11 12 @Override 13 public TaskPO saveOrUpdate(TaskVO taskVO) { 14 return taskService.save(input); 15 } 16 17 18 @Override 19 public Page<TaskPO> hashList(TaskReqVO taskReqVO) { 20 Page<TaskPO> page = Page.of(taskReqVO.getIndex(), taskReqVO.getSize()); 21 LambdaQueryWrapper<TaskPO> wrapper = Wrappers.<TaskPO>lambdaQuery() 22 .in(CollectionUtils.isNotEmpty(taskReqVO.getStatus()), TaskPO::getStatus, taskReqVO.getStatus()) 23 .eq(taskReqVO.getBizType() != null, TaskPO::getBizType, taskReqVO.getBizType()) 24 .le(taskReqVO.getErrCnt() != null, TaskPO::getErrCnt, taskReqVO.getErrCnt()) 25 .apply("mod(id," + taskReqVO.getShardTotal() + ") =" + taskReqVO.getShardIndex() + " ") 26 .orderByAsc(TaskPO::getCreateTime); 27 return taskService.page(page, wrapper); 28 } 29 30 private TaskVO getTask(String fileName,String pin, String key,TaskBizTypeEnum bizType,TaskCreateTypeEnum taskCreateType){ 31 // build task 32 return res; 33 } 34 35 @Override 36 public TaskDTO pullTask(){ 37 TaskDTO target = null; 38 ShardingUtil.ShardingVO shardingVo = ShardingUtil.getShardingVo(); 39 log.info("iscWorker.pullTask workerIndex: {}, total:{}" , shardingVo.getIndex(),shardingVo.getTotal()); 40 TaskReqVO queryDTO = new TaskReqVO(); 41 queryDTO.setShardIndex(shardingVo.getIndex()); 42 queryDTO.setShardTotal(shardingVo.getTotal()); 43 queryDTO.setStatus(Lists.newArrayList(TaskStatusEnum.PENDING,TaskStatusEnum.FAILURE)); 44 queryDTO.setErrCnt(MAX_ERR_CNT); 45 queryDTO.setIndex(0L); 46 queryDTO.setSize(LIMIT); 47 Page<TaskPO> targetList = hashList(queryDTO); 48 if(CollectionUtils.isNotEmpty(targetList.getRecords())){ 49 log.info("PublishMkuBySkuWorker.pullTask 准备执行:{}" , JSON.toJSONString(targetList)); 50 target = new TaskDTO<>(targetList.getRecords().get(0)); 51 } 52 return target; 53 } 54 55 56 @Override 57 public Boolean error(TaskVO taskInfo) { 58 return task.update(taskInfo); 59 } 60 61 /** 失败处理*/ 62 @Override 63 public void errHandle(TaskDTO task, String errMsg){ 64 error(new TaskVO(task.getTaskId())); 65 Profiler.businessAlarm(UmpKeyConstant.BUSINESS_KEY_TASK_WARNING,("excel批量导入-任务执行异常:"+errMsg+task.getTaskId())); 66 log.info("TaskMangeServiceImpl.errHandle 任务Id{}执行失败:{}", task.getTaskId(),errMsg); 67 } 68 69} 70 71 72
五、类图

图2 类图
作者:京东工业 于洋
来源:京东云开发者社区 转载请注明来源
