DatahubJobTask.java 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167
  1. package com.mokamrp.privates.task;
  2. import com.aliyun.datahub.client.DatahubClient;
  3. import com.aliyun.datahub.client.exception.DatahubClientException;
  4. import com.aliyun.datahub.client.model.PutRecordsResult;
  5. import com.aliyun.datahub.client.model.RecordEntry;
  6. import com.aliyun.datahub.client.model.RecordSchema;
  7. import com.aliyun.datahub.client.model.TupleRecordData;
  8. import com.mokamrp.privates.entity.JobContactEmployee;
  9. import com.mokamrp.privates.mapper.vo.WorkContactEmployeeHubVo;
  10. import com.mokamrp.privates.service.WorkContactEmployeeService;
  11. import com.mokamrp.privates.utils.RedisUtil;
  12. import net.javacrumbs.shedlock.core.SchedulerLock;
  13. import org.springframework.beans.factory.annotation.Autowired;
  14. import org.springframework.beans.factory.annotation.Value;
  15. import org.springframework.scheduling.annotation.Scheduled;
  16. import org.springframework.stereotype.Component;
  17. import javax.annotation.Resource;
  18. import java.text.SimpleDateFormat;
  19. import java.util.ArrayList;
  20. import java.util.Date;
  21. import java.util.List;
  22. @Component
  23. public class DatahubJobTask {
  24. @Autowired
  25. private WorkContactEmployeeService workContactEmployeeService;
  26. @Resource
  27. private RedisUtil redisUtil;
  28. @Resource
  29. private DatahubClient datahubClient;
  30. @Value("${datahub.topic}")
  31. private String topic;
  32. @Value("${datahub.project}")
  33. private String project;
  34. /**
  35. * 导出任务,单进程
  36. */
  37. @Scheduled(fixedDelay = 30000)
  38. @SchedulerLock(name = "pangu:DatahubJobTask", lockAtLeastForString = "PT5S", lockAtMostForString = "PT15M")
  39. public void execute() {
  40. System.out.println("执行datahub任务,当前时间"+(new Date())+",当前长度"+redisUtil.lGetListSize(JobContactEmployee.QUEUE_KEY));
  41. System.out.println(topic);
  42. System.out.println(project);
  43. while (redisUtil.lGetListSize(JobContactEmployee.QUEUE_KEY) >=1) {
  44. JobContactEmployee out = (JobContactEmployee) redisUtil.rightPop(JobContactEmployee.QUEUE_KEY);
  45. job(project,topic,3,out.getCorpId());
  46. break;
  47. }
  48. }
  49. public void job(String project,String topic,int retryTimes,Integer corpId){
  50. List<WorkContactEmployeeHubVo> list;
  51. Integer maxId;
  52. maxId = null;
  53. while(true){
  54. list = workContactEmployeeService.getCorpEmployeeContact(corpId,maxId);
  55. if(list == null || list.size() <=0){
  56. System.out.println("执行完毕CorpId"+corpId);
  57. break;
  58. }
  59. maxId = list.get(list.size()-1).getId();
  60. tupleExample(project,topic,retryTimes,list);
  61. }
  62. }
  63. // 写入Tuple型数据
  64. public void tupleExample(String project,String topic,int retryTimes, List<WorkContactEmployeeHubVo> list) {
  65. Date date = new Date();
  66. String timestamp = String.valueOf((new Date()).getTime()/1000);
  67. System.out.println(timestamp);
  68. Integer time = Integer.valueOf(timestamp);
  69. String strDateFormat = "yyyyMMdd";
  70. SimpleDateFormat sdf = new SimpleDateFormat(strDateFormat);
  71. String resDate = sdf.format(new Date());
  72. WorkContactEmployeeHubVo vo;
  73. // 获取schema
  74. RecordSchema recordSchema = datahubClient.getTopic(project,topic ).getRecordSchema();
  75. // 推送数据,生成十条数据
  76. Integer createtime;
  77. List<RecordEntry> recordEntries = new ArrayList<>();
  78. for (int i = 0; i < list.size(); i++) {
  79. createtime = 0;
  80. RecordEntry recordEntry = new RecordEntry();
  81. // 对每条数据设置额外属性,例如ip 机器名等。可以不设置额外属性,不影响数据写入
  82. //recordEntry.addAttribute("key1", "value1");
  83. TupleRecordData data = new TupleRecordData(recordSchema);
  84. vo = list.get(i);
  85. if(vo.getFollowCreateTime() != null && !vo.getFollowCreateTime().isEmpty()){
  86. createtime = date2TimeStamp(vo.getFollowCreateTime(),"yyyy-MM-dd HH:mm:ss");
  87. }
  88. data.setField("external_userid", vo.getExternalUserid());
  89. data.setField("external_name", vo.getExternalName());
  90. data.setField("external_type", vo.getExternalType());
  91. data.setField("external_avatar", vo.getExternalAvatar());
  92. data.setField("external_corp_name", vo.getExternalCorpName());
  93. data.setField("external_corp_full_name", vo.getExternalCorpFullName());
  94. data.setField("external_gender", vo.getExternalGender());
  95. data.setField("wx_corpid", vo.getWxCorpid());
  96. data.setField("follow_userid", vo.getFollowUserid());
  97. data.setField("follow_remark", vo.getFollowRemark());
  98. data.setField("follow_createtime", createtime);
  99. data.setField("follow_remark_mobiles", vo.getFollowRemark());
  100. data.setField("follow_add_way", vo.getFollowAddWay());
  101. data.setField("follow_tag_id", "");
  102. data.setField("report_time", time);
  103. data.setField("ds", resDate);
  104. recordEntry.setRecordData(data);
  105. recordEntries.add(recordEntry);
  106. }
  107. try {
  108. PutRecordsResult result = datahubClient.putRecords(project, topic, recordEntries);
  109. int i = result.getFailedRecordCount();
  110. if (i > 0) {
  111. retry(datahubClient, result.getFailedRecords(), retryTimes, project, topic);
  112. }
  113. } catch (DatahubClientException e) {
  114. System.out.println("requestId:" + e.getRequestId() + "\tmessage:" + e.getErrorMessage());
  115. }
  116. return ;
  117. }
  118. //重试机制
  119. public static void retry(DatahubClient client, List<RecordEntry> records, int retryTimes, String project, String topic) {
  120. boolean suc = false;
  121. while (retryTimes != 0) {
  122. retryTimes = retryTimes - 1;
  123. PutRecordsResult recordsResult = client.putRecords(project, topic, records);
  124. if (recordsResult.getFailedRecordCount() > 0) {
  125. retry(client,recordsResult.getFailedRecords(),retryTimes,project,topic);
  126. }
  127. suc = true;
  128. break;
  129. }
  130. if (!suc) {
  131. System.out.println("retryFailure");
  132. }
  133. }
  134. public static Integer date2TimeStamp(String date_str,String format){
  135. try {
  136. SimpleDateFormat sdf = new SimpleDateFormat(format);
  137. return Integer.valueOf(String.valueOf(sdf.parse(date_str).getTime()/1000));
  138. } catch (Exception e) {
  139. e.printStackTrace();
  140. }
  141. return 0;
  142. }
  143. }