| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167 |
- 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<WorkContactEmployeeHubVo> 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<WorkContactEmployeeHubVo> 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<RecordEntry> 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<RecordEntry> 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;
- }
- }
|