package com.mokamrp.privates.task; import com.aliyun.datahub.client.DatahubClient; import com.aliyun.datahub.client.exception.DatahubClientException; import com.aliyun.datahub.client.model.PutRecordsResult; import com.aliyun.datahub.client.model.RecordEntry; import com.aliyun.datahub.client.model.RecordSchema; import com.aliyun.datahub.client.model.TupleRecordData; import com.mokamrp.privates.entity.JobContactEmployee; import com.mokamrp.privates.mapper.vo.WorkContactEmployeeHubVo; import com.mokamrp.privates.service.WorkContactEmployeeService; import com.mokamrp.privates.utils.RedisUtil; import net.javacrumbs.shedlock.core.SchedulerLock; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import javax.annotation.Resource; import java.text.SimpleDateFormat; import java.util.ArrayList; import java.util.Date; import java.util.List; @Component public class DatahubJobTask { @Autowired private WorkContactEmployeeService workContactEmployeeService; @Resource private RedisUtil redisUtil; @Resource private DatahubClient datahubClient; @Value("${datahub.topic}") private String topic; @Value("${datahub.project}") private String project; /** * 导出任务,单进程 */ @Scheduled(fixedDelay = 30000) @SchedulerLock(name = "pangu:DatahubJobTask", lockAtLeastForString = "PT5S", lockAtMostForString = "PT15M") public void execute() { System.out.println("执行datahub任务,当前时间"+(new Date())+",当前长度"+redisUtil.lGetListSize(JobContactEmployee.QUEUE_KEY)); System.out.println(topic); System.out.println(project); while (redisUtil.lGetListSize(JobContactEmployee.QUEUE_KEY) >=1) { JobContactEmployee out = (JobContactEmployee) redisUtil.rightPop(JobContactEmployee.QUEUE_KEY); job(project,topic,3,out.getCorpId()); break; } } public void job(String project,String topic,int retryTimes,Integer corpId){ List list; Integer maxId; maxId = null; while(true){ list = workContactEmployeeService.getCorpEmployeeContact(corpId,maxId); if(list == null || list.size() <=0){ System.out.println("执行完毕CorpId"+corpId); break; } maxId = list.get(list.size()-1).getId(); tupleExample(project,topic,retryTimes,list); } } // 写入Tuple型数据 public void tupleExample(String project,String topic,int retryTimes, List list) { Date date = new Date(); String timestamp = String.valueOf((new Date()).getTime()/1000); System.out.println(timestamp); Integer time = Integer.valueOf(timestamp); String strDateFormat = "yyyyMMdd"; SimpleDateFormat sdf = new SimpleDateFormat(strDateFormat); String resDate = sdf.format(new Date()); WorkContactEmployeeHubVo vo; // 获取schema RecordSchema recordSchema = datahubClient.getTopic(project,topic ).getRecordSchema(); // 推送数据,生成十条数据 Integer createtime; List recordEntries = new ArrayList<>(); for (int i = 0; i < list.size(); i++) { createtime = 0; RecordEntry recordEntry = new RecordEntry(); // 对每条数据设置额外属性,例如ip 机器名等。可以不设置额外属性,不影响数据写入 //recordEntry.addAttribute("key1", "value1"); TupleRecordData data = new TupleRecordData(recordSchema); vo = list.get(i); if(vo.getFollowCreateTime() != null && !vo.getFollowCreateTime().isEmpty()){ createtime = date2TimeStamp(vo.getFollowCreateTime(),"yyyy-MM-dd HH:mm:ss"); } data.setField("external_userid", vo.getExternalUserid()); data.setField("external_name", vo.getExternalName()); data.setField("external_type", vo.getExternalType()); data.setField("external_avatar", vo.getExternalAvatar()); data.setField("external_corp_name", vo.getExternalCorpName()); data.setField("external_corp_full_name", vo.getExternalCorpFullName()); data.setField("external_gender", vo.getExternalGender()); data.setField("wx_corpid", vo.getWxCorpid()); data.setField("follow_userid", vo.getFollowUserid()); data.setField("follow_remark", vo.getFollowRemark()); data.setField("follow_createtime", createtime); data.setField("follow_remark_mobiles", vo.getFollowRemark()); data.setField("follow_add_way", vo.getFollowAddWay()); data.setField("follow_tag_id", ""); data.setField("report_time", time); data.setField("ds", resDate); recordEntry.setRecordData(data); recordEntries.add(recordEntry); } try { PutRecordsResult result = datahubClient.putRecords(project, topic, recordEntries); int i = result.getFailedRecordCount(); if (i > 0) { retry(datahubClient, result.getFailedRecords(), retryTimes, project, topic); } } catch (DatahubClientException e) { System.out.println("requestId:" + e.getRequestId() + "\tmessage:" + e.getErrorMessage()); } return ; } //重试机制 public static void retry(DatahubClient client, List records, int retryTimes, String project, String topic) { boolean suc = false; while (retryTimes != 0) { retryTimes = retryTimes - 1; PutRecordsResult recordsResult = client.putRecords(project, topic, records); if (recordsResult.getFailedRecordCount() > 0) { retry(client,recordsResult.getFailedRecords(),retryTimes,project,topic); } suc = true; break; } if (!suc) { System.out.println("retryFailure"); } } public static Integer date2TimeStamp(String date_str,String format){ try { SimpleDateFormat sdf = new SimpleDateFormat(format); return Integer.valueOf(String.valueOf(sdf.parse(date_str).getTime()/1000)); } catch (Exception e) { e.printStackTrace(); } return 0; } }