JobServiceKafka.java 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108
  1. package com.java110.job.kafka;
  2. import com.alibaba.fastjson.JSONObject;
  3. import com.java110.core.base.controller.BaseController;
  4. import com.java110.core.context.BusinessServiceDataFlow;
  5. import com.java110.core.factory.DataTransactionFactory;
  6. import com.java110.job.smo.IJobServiceSMO;
  7. import com.java110.utils.constant.KafkaConstant;
  8. import com.java110.utils.constant.ResponseConstant;
  9. import com.java110.utils.constant.StatusConstant;
  10. import com.java110.utils.exception.InitConfigDataException;
  11. import com.java110.utils.exception.InitDataFlowContextException;
  12. import com.java110.utils.kafka.KafkaFactory;
  13. import com.java110.utils.util.Assert;
  14. import org.apache.kafka.clients.consumer.ConsumerRecord;
  15. import org.slf4j.Logger;
  16. import org.slf4j.LoggerFactory;
  17. import org.springframework.beans.factory.annotation.Autowired;
  18. import org.springframework.kafka.annotation.KafkaListener;
  19. import java.util.HashMap;
  20. import java.util.Map;
  21. /**
  22. * kafka侦听
  23. * Created by wuxw on 2018/4/15.
  24. */
  25. public class JobServiceKafka extends BaseController {
  26. private final static Logger logger = LoggerFactory.getLogger(JobServiceKafka.class);
  27. @Autowired
  28. private IJobServiceSMO jobServiceSMOImpl;
  29. @KafkaListener(topics = {"jobServiceTopic"})
  30. public void listen(ConsumerRecord<?, ?> record) {
  31. logger.info("kafka的key: " + record.key());
  32. logger.info("kafka的value: " + record.value().toString());
  33. String orderInfo = record.value().toString();
  34. BusinessServiceDataFlow businessServiceDataFlow = null;
  35. JSONObject responseJson = null;
  36. try {
  37. Map<String, String> headers = new HashMap<String, String>();
  38. //预校验
  39. preValiateOrderInfo(orderInfo);
  40. businessServiceDataFlow = this.writeDataToDataFlowContext(orderInfo, headers);
  41. //responseJson = jobServiceSMOImpl.service(businessServiceDataFlow);
  42. } catch (InitDataFlowContextException e) {
  43. logger.error("请求报文错误,初始化 BusinessServiceDataFlow失败" + orderInfo, e);
  44. responseJson = DataTransactionFactory.createNoBusinessTypeBusinessResponseJson(orderInfo, ResponseConstant.RESULT_PARAM_ERROR, e.getMessage(), null);
  45. } catch (InitConfigDataException e) {
  46. logger.error("请求报文错误,加载配置信息失败" + orderInfo, e);
  47. responseJson = DataTransactionFactory.createNoBusinessTypeBusinessResponseJson(orderInfo, ResponseConstant.RESULT_PARAM_ERROR, e.getMessage(), null);
  48. } catch (Exception e) {
  49. logger.error("请求订单异常", e);
  50. responseJson = DataTransactionFactory.createBusinessResponseJson(businessServiceDataFlow, ResponseConstant.RESULT_CODE_ERROR, e.getMessage() + e,
  51. null);
  52. } finally {
  53. logger.debug("当前请求报文:" + orderInfo + ", 当前返回报文:" + responseJson.toJSONString());
  54. //只有business 和 instance 过程才做通知消息
  55. if (!StatusConstant.REQUEST_BUSINESS_TYPE_BUSINESS.equals(responseJson.getString("businessType"))
  56. && !StatusConstant.REQUEST_BUSINESS_TYPE_INSTANCE.equals(responseJson.getString("businessType"))) {
  57. return;
  58. }
  59. try {
  60. KafkaFactory.sendKafkaMessage(KafkaConstant.TOPIC_NOTIFY_CENTER_SERVICE_NAME, "", responseJson.toJSONString());
  61. } catch (Exception e) {
  62. logger.error("用户服务通知centerService失败" + responseJson, e);
  63. //这里保存异常信息
  64. }
  65. }
  66. }
  67. @KafkaListener(topics = {"${kafka.hcGovTopic}"})
  68. public void hcGovListen(ConsumerRecord<?, ?> record) {
  69. logger.info("kafka的key: " + record.key());
  70. logger.info("kafka的value: " + record.value().toString());
  71. String orderInfo = record.value().toString();
  72. try {
  73. logger.debug("hcGovkafka 接收到数据", orderInfo);
  74. //responseJson = jobServiceSMOImpl.service(businessServiceDataFlow);
  75. } catch (Exception e) {
  76. logger.error("请求订单异常", e);
  77. } finally {
  78. }
  79. }
  80. /**
  81. * 这里预校验,请求报文中不能有 dataFlowId
  82. *
  83. * @param orderInfo
  84. */
  85. private void preValiateOrderInfo(String orderInfo) {
  86. JSONObject reqJson = JSONObject.parseObject(orderInfo);
  87. Assert.hasKeyAndValue(reqJson, "header", "请求报文中未包含header");
  88. Assert.hasKeyAndValue(reqJson, "body", "请求报文中未包含body");
  89. JSONObject header = reqJson.getJSONObject("header");
  90. Assert.hasKeyAndValue(header, "serviceCode", "请求报文中未包含serviceCode");
  91. Assert.hasKeyAndValue(header, "sign", "请求报文中未包含sign");
  92. Assert.hasKeyAndValue(header, "resTime", "请求报文中未包含reqTime");
  93. Assert.hasKeyAndValue(header, "code", "请求报文中未包含reqTime");
  94. Assert.hasKeyAndValue(header, "msg", "请求报文中未包含reqTime");
  95. }
  96. }