通过mq消息同步签名签章

feature-LeeD-20221025
liuliang 2 years ago
parent baf973d41d
commit 3e37c9c204

@ -0,0 +1,22 @@
package weaver.interfaces.dito.mq;
import com.alibaba.fastjson.JSONArray;
public interface ConsumerBase {
/**
*
* ele_signature
* @param tableName ele_signature
* @return
*/
boolean support(String tableName);
/**
*
* @param jsonArray
* @param tableName
*/
void consumer(JSONArray jsonArray, String tableName);
}

@ -1,40 +1,40 @@
//package weaver.interfaces.dito.mq; package weaver.interfaces.dito.mq;
//
//import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
//import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
//import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
//import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
//import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.message.MessageExt;
//import weaver.general.BaseBean; import weaver.general.BaseBean;
//
//import java.io.UnsupportedEncodingException; import java.io.UnsupportedEncodingException;
//import java.util.List; import java.util.List;
//
//
//public class HrmRocketMsgListener implements MessageListenerConcurrently { public class HrmRocketMsgListener implements MessageListenerConcurrently {
// @Override @Override
// public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext consumeConcurrentlyContext) { public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
// BaseBean bb = new BaseBean(); BaseBean bb = new BaseBean();
// MessageExt msg = msgs.get(0); MessageExt msg = msgs.get(0);
//
// HrmRocketmqUtil hrmRocketmqUtil = new HrmRocketmqUtil(); HrmRocketmqUtil hrmRocketmqUtil = new HrmRocketmqUtil();
//
// try { try {
//
// bb.writeLog("Consumer---6----"+new String(msg.getBody())); bb.writeLog("Consumer---6----"+new String(msg.getBody()));
// String msgdata = new String(msg.getBody(),"UTF-8"); String msgdata = new String(msg.getBody(),"UTF-8");
// if(StringUtils.isBlank(msgdata)) if(StringUtils.isBlank(msgdata))
// { {
// int errcount = hrmRocketmqUtil.updateOrgData(msgdata); int errcount = hrmRocketmqUtil.updateOrgData(msgdata);
// bb.writeLog("Consumer---errcount---:"+errcount); bb.writeLog("Consumer---errcount---:"+errcount);
// } }
// } catch (UnsupportedEncodingException e) { } catch (UnsupportedEncodingException e) {
// e.printStackTrace(); e.printStackTrace();
// bb.writeLog("Consumer---UnsupportedEncodingException---e:"+e); bb.writeLog("Consumer---UnsupportedEncodingException---e:"+e);
// return ConsumeConcurrentlyStatus.RECONSUME_LATER; return ConsumeConcurrentlyStatus.RECONSUME_LATER;
// } }
// return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
//
// } }
//
//} }

@ -1,182 +1,182 @@
//package weaver.interfaces.dito.mq; package weaver.interfaces.dito.mq;
//
//
//import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
//import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.exception.MQClientException;
//import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
//import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
//import weaver.file.Prop; import weaver.file.Prop;
//import weaver.general.BaseBean; import weaver.general.BaseBean;
//import weaver.general.GCONST; import weaver.general.GCONST;
//import weaver.general.InitServer; import weaver.general.InitServer;
//import weaver.interfaces.dito.comInfo.PropBean; import weaver.interfaces.dito.comInfo.PropBean;
//
//import javax.servlet.ServletException; import javax.servlet.ServletException;
//import javax.servlet.http.HttpServlet; import javax.servlet.http.HttpServlet;
//import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletRequest;
//import javax.servlet.http.HttpServletResponse; import javax.servlet.http.HttpServletResponse;
//import java.io.IOException; import java.io.IOException;
//import java.text.SimpleDateFormat; import java.text.SimpleDateFormat;
//import java.util.ArrayList; import java.util.ArrayList;
//import java.util.Date; import java.util.Date;
//
//
//public class HrmRocketmqServlet extends HttpServlet { public class HrmRocketmqServlet extends HttpServlet {
//
//
// private boolean isMainIp = false;//是否为主节点 private boolean isMainIp = false;//是否为主节点
//
// private boolean isSampleMode = false;//是否为单机启用ridis并且配置了主节点的则不是单机否则认为是集群环境 private boolean isSampleMode = false;//是否为单机启用ridis并且配置了主节点的则不是单机否则认为是集群环境
//
// @Override @Override
// public void init() throws ServletException public void init() throws ServletException
// { {
//
// BaseBean bb = new BaseBean(); BaseBean bb = new BaseBean();
// SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmmss"); SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmmss");
//
// bb.writeLog("initiated---进入时间:"+sdf.format(new Date())); bb.writeLog("initiated---进入时间:"+sdf.format(new Date()));
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
//
// try{ try{
// isMainIp = isMainIp();//计算是否为集群中的主节点 isMainIp = isMainIp();//计算是否为集群中的主节点
// isSampleMode = isSampleMode();//判断是否为单机 isSampleMode = isSampleMode();//判断是否为单机
// bb.writeLog("isMainIp:"+isMainIp); bb.writeLog("isMainIp:"+isMainIp);
// bb.writeLog("isSampleMode:"+isSampleMode); bb.writeLog("isSampleMode:"+isSampleMode);
//
// if(isSampleMode){ if(isSampleMode){
// initData(); initData();
// }else{ }else{
// if (isMainIp) { if (isMainIp) {
// //则调用service方法加载设置重新获取workflowid相关超时设置OvertimeEntity //则调用service方法加载设置重新获取workflowid相关超时设置OvertimeEntity
// initData(); initData();
// } }
// } }
//
// }catch (Exception e){ }catch (Exception e){
// bb.writeLog("Consumer resource model initiated--MQClientException:"+e); bb.writeLog("Consumer resource model initiated--MQClientException:"+e);
// } }
//
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// } }
//
//
// public void initData(){ public void initData(){
// BaseBean bb = new BaseBean(); BaseBean bb = new BaseBean();
// PropBean propBean = new PropBean(); PropBean propBean = new PropBean();
//
// try { try {
//// String consumerGroup = propBean.getUfPropValueStatic("consumerGroup"); // String consumerGroup = propBean.getUfPropValueStatic("consumerGroup");
//// bb.writeLog("consumerGroup:" + consumerGroup); // bb.writeLog("consumerGroup:" + consumerGroup);
// String hrmConsumerGroup = PropBean.getUfPropValue("hrmConsumerGroup"); String hrmConsumerGroup = PropBean.getUfPropValue("hrmConsumerGroup");
// String hrmConsumerAddr = PropBean.getUfPropValue("hrmConsumerAddr"); String hrmConsumerAddr = PropBean.getUfPropValue("hrmConsumerAddr");
// String hrmInstanceName = PropBean.getUfPropValue("hrmInstanceName"); String hrmInstanceName = PropBean.getUfPropValue("hrmInstanceName");
// String hrmSubExpr = PropBean.getUfPropValue("hrmSubExpr"); String hrmSubExpr = PropBean.getUfPropValue("hrmSubExpr");
// String hrmMQAuthID = PropBean.getUfPropValue("hrmMQAuthID"); String hrmMQAuthID = PropBean.getUfPropValue("hrmMQAuthID");
// String hrmMQAuthPWD = PropBean.getUfPropValue("hrmMQAuthPWD"); String hrmMQAuthPWD = PropBean.getUfPropValue("hrmMQAuthPWD");
// String hrmMQClusterName = PropBean.getUfPropValue("hrmMQClusterName"); String hrmMQClusterName = PropBean.getUfPropValue("hrmMQClusterName");
// String hrmMQTenantID = PropBean.getUfPropValue("hrmMQTenantID"); String hrmMQTenantID = PropBean.getUfPropValue("hrmMQTenantID");
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(hrmConsumerGroup); DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(hrmConsumerGroup);
//// String namesrvAddr = propBean.getUfPropValueStatic("namesrvAddr"); // String namesrvAddr = propBean.getUfPropValueStatic("namesrvAddr");
//// bb.writeLog("namesrvAddr:" + namesrvAddr); // bb.writeLog("namesrvAddr:" + namesrvAddr);
// consumer.setNamesrvAddr(hrmConsumerAddr); consumer.setNamesrvAddr(hrmConsumerAddr);
//// String instanceName = propBean.getUfPropValueStatic("instanceName"); // String instanceName = propBean.getUfPropValueStatic("instanceName");
//// bb.writeLog("instanceName:" + instanceName); // bb.writeLog("instanceName:" + instanceName);
// consumer.setInstanceName(hrmInstanceName); consumer.setInstanceName(hrmInstanceName);
// consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
//// String topic = propBean.getUfPropValueStatic("topic"); // String topic = propBean.getUfPropValueStatic("topic");
//// String subExpression = propBean.getUfPropValueStatic("subExpression"); // String subExpression = propBean.getUfPropValueStatic("subExpression");
//// bb.writeLog("topic:" + topic); // bb.writeLog("topic:" + topic);
//// bb.writeLog("subExpression:" + subExpression); // bb.writeLog("subExpression:" + subExpression);
//// String authid = propBean.getUfPropValueStatic("authid"); // String authid = propBean.getUfPropValueStatic("authid");
//// String authpwd = propBean.getUfPropValueStatic("authpwd"); // String authpwd = propBean.getUfPropValueStatic("authpwd");
//// String clustername = propBean.getUfPropValueStatic("clustername"); // String clustername = propBean.getUfPropValueStatic("clustername");
//// String tenantid = propBean.getUfPropValueStatic("tenantid"); // String tenantid = propBean.getUfPropValueStatic("tenantid");
//// bb.writeLog("authid:" + authid); // bb.writeLog("authid:" + authid);
//// bb.writeLog("authpwd:" + authpwd); // bb.writeLog("authpwd:" + authpwd);
//// bb.writeLog("clustername:" + clustername); // bb.writeLog("clustername:" + clustername);
//// bb.writeLog("tenantid:" + tenantid); // bb.writeLog("tenantid:" + tenantid);
// consumer.subscribe(hrmInstanceName, hrmSubExpr); consumer.subscribe(hrmInstanceName, hrmSubExpr);
// consumer.setAuthID(hrmMQAuthID); consumer.setAuthID(hrmMQAuthID);
// consumer.setAuthPWD(hrmMQAuthPWD); consumer.setAuthPWD(hrmMQAuthPWD);
// consumer.setClusterName(hrmMQClusterName); consumer.setClusterName(hrmMQClusterName);
// consumer.setTenantID(hrmMQTenantID); consumer.setTenantID(hrmMQTenantID);
//
// consumer.setConsumeThreadMin(1); consumer.setConsumeThreadMin(1);
// consumer.setConsumeThreadMax(1); consumer.setConsumeThreadMax(1);
// consumer.setConsumeMessageBatchMaxSize(1); consumer.setConsumeMessageBatchMaxSize(1);
//
// consumer.setMessageModel(MessageModel.BROADCASTING); consumer.setMessageModel(MessageModel.BROADCASTING);
// consumer.registerMessageListener(new HrmRocketMsgListener()); consumer.registerMessageListener(new HrmRocketMsgListener());
// consumer.start(); consumer.start();
// bb.writeLog("Consumer88Started."); bb.writeLog("Consumer88Started.");
//
// } catch (MQClientException var13) { } catch (MQClientException var13) {
// bb.writeLog("Consumer resource model initiated--MQClientException:" + var13); bb.writeLog("Consumer resource model initiated--MQClientException:" + var13);
// } }
// } }
//
// /** /**
// * 实现 HttpServlet 的 doGet 方法,不作任何操作 * HttpServlet doGet
// * @param request * @param request
// * @param response * @param response
// * @throws ServletException * @throws ServletException
// * @throws IOException * @throws IOException
// */ */
// @Override @Override
// public void doGet(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException { public void doGet(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {
// doPost(request, response); doPost(request, response);
// } }
//
// /** /**
// * 实现 HttpServlet 的 doPost 方法,不作任何操作 * HttpServlet doPost
// * @param request * @param request
// * @param response * @param response
// * @throws ServletException * @throws ServletException
// * @throws IOException * @throws IOException
// */ */
//
// @Override @Override
// public void doPost(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException { public void doPost(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {
// } }
//
// @Override @Override
// public void destroy() public void destroy()
// { {
// // 什么也不做 // 什么也不做
// } }
//
//
// private boolean isMainIp() { private boolean isMainIp() {
// new BaseBean().writeLog("超时判断主节点"); new BaseBean().writeLog("超时判断主节点");
// String mainControlIp = ""; String mainControlIp = "";
// ArrayList<String> hostIps = new InitServer().getRealIp(); ArrayList<String> hostIps = new InitServer().getRealIp();
// Prop prop = Prop.getInstance(); Prop prop = Prop.getInstance();
// mainControlIp = prop.getPropValue(GCONST.getConfigFile(), "MainControlIP"); mainControlIp = prop.getPropValue(GCONST.getConfigFile(), "MainControlIP");
// if (hostIps == null || hostIps.size() == 0) { if (hostIps == null || hostIps.size() == 0) {
// new BaseBean().writeLog("System Init Error:Cannot get local Ip address,This may cause scripts or Timed task not run! "); new BaseBean().writeLog("System Init Error:Cannot get local Ip address,This may cause scripts or Timed task not run! ");
// } else { } else {
// new BaseBean().writeLog("System Init Message:mainControlIp=" + mainControlIp + " localIp:" + hostIps.toString()); new BaseBean().writeLog("System Init Message:mainControlIp=" + mainControlIp + " localIp:" + hostIps.toString());
// } }
// if ((!"".equals(mainControlIp) && hostIps.contains(mainControlIp)) || "".equals(mainControlIp)) { if ((!"".equals(mainControlIp) && hostIps.contains(mainControlIp)) || "".equals(mainControlIp)) {
// return true; return true;
// } }
// return false; return false;
// } }
//
// public boolean isSampleMode() { public boolean isSampleMode() {
// BaseBean base = new BaseBean(); BaseBean base = new BaseBean();
// boolean redis_flag = "1".equals(base.getPropValue("weaver_new_session", "status")); boolean redis_flag = "1".equals(base.getPropValue("weaver_new_session", "status"));
// boolean mainIp_flag = !"".equals(base.getPropValue("weaver", "MainControlIP")); boolean mainIp_flag = !"".equals(base.getPropValue("weaver", "MainControlIP"));
// base.writeLog("超时判断是否为集群环境redis_flag======"+redis_flag); base.writeLog("超时判断是否为集群环境redis_flag======"+redis_flag);
// base.writeLog("超时判断是否为集群环境mainIp_flag======="+mainIp_flag); base.writeLog("超时判断是否为集群环境mainIp_flag======="+mainIp_flag);
// return !(redis_flag && mainIp_flag); return !(redis_flag && mainIp_flag);
// } }
//
//
//
//} }

@ -1,40 +1,41 @@
//package weaver.interfaces.dito.mq; package weaver.interfaces.dito.mq;
//
//import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
//import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
//import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
//import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.message.MessageExt;
//import weaver.general.BaseBean; import weaver.general.BaseBean;
//
//import java.io.UnsupportedEncodingException; import java.io.UnsupportedEncodingException;
//import java.util.List; import java.util.List;
//
//
//public class RocketMsgListener implements MessageListenerConcurrently { public class RocketMsgListener implements MessageListenerConcurrently {
// @Override @Override
// public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext consumeConcurrentlyContext) { public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
// BaseBean bb = new BaseBean(); BaseBean bb = new BaseBean();
// MessageExt msg = msgs.get(0); MessageExt msg = msgs.get(0);
//
// RocketmqUtil rocketmqUtil = new RocketmqUtil(); RocketmqUtil rocketmqUtil = new RocketmqUtil();
//
// try { try {
//
// bb.writeLog("Consumer---3----"+new String(msg.getBody())); bb.writeLog("Consumer---3----"+new String(msg.getBody()));
// String msgdata = new String(msg.getBody(),"UTF-8");
// if(!"".equals(msgdata)) String msgdata = new String(msg.getBody(),"UTF-8");
// { if(!"".equals(msgdata))
// String data = msgdata.substring(msgdata.indexOf("{")); {
// int errcount = rocketmqUtil.updateOrgData(data); String data = msgdata.substring(msgdata.indexOf("{"));
// bb.writeLog("Consumer---errcount---:"+errcount); int errcount = rocketmqUtil.updateOrgData(data);
// } bb.writeLog("Consumer---errcount---:"+errcount);
// } catch (UnsupportedEncodingException e) { }
// e.printStackTrace(); } catch (UnsupportedEncodingException e) {
// bb.writeLog("Consumer---UnsupportedEncodingException---e:"+e); e.printStackTrace();
// return ConsumeConcurrentlyStatus.RECONSUME_LATER; bb.writeLog("Consumer---UnsupportedEncodingException---e:"+e);
// } return ConsumeConcurrentlyStatus.RECONSUME_LATER;
// return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }
// return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
// }
// }
//}
}

@ -1,216 +1,216 @@
//package weaver.interfaces.dito.mq; package weaver.interfaces.dito.mq;
//
//
//import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
//import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.exception.MQClientException;
//import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
//import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
//import weaver.file.Prop; import weaver.file.Prop;
//import weaver.general.BaseBean; import weaver.general.BaseBean;
//import weaver.general.GCONST; import weaver.general.GCONST;
//import weaver.general.InitServer; import weaver.general.InitServer;
//import weaver.interfaces.dito.comInfo.PropBean; import weaver.interfaces.dito.comInfo.PropBean;
//
//import javax.servlet.ServletException; import javax.servlet.ServletException;
//import javax.servlet.http.HttpServlet; import javax.servlet.http.HttpServlet;
//import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletRequest;
//import javax.servlet.http.HttpServletResponse; import javax.servlet.http.HttpServletResponse;
//import java.io.IOException; import java.io.IOException;
//import java.text.SimpleDateFormat; import java.text.SimpleDateFormat;
//import java.util.ArrayList; import java.util.ArrayList;
//import java.util.Date; import java.util.Date;
//
//
//
//public class RocketmqServlet extends HttpServlet { public class RocketmqServlet extends HttpServlet {
//
//
// private boolean isMainIp = false;//是否为主节点 private boolean isMainIp = false;//是否为主节点
//
// private boolean isSampleMode = false;//是否为单机启用ridis并且配置了主节点的则不是单机否则认为是集群环境 private boolean isSampleMode = false;//是否为单机启用ridis并且配置了主节点的则不是单机否则认为是集群环境
//
// @Override @Override
// public void init() throws ServletException public void init() throws ServletException
// { {
//
// BaseBean bb = new BaseBean(); BaseBean bb = new BaseBean();
// SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmmss"); SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmmss");
//
// bb.writeLog("initiated---进入时间:"+sdf.format(new Date())); bb.writeLog("initiated---进入时间:"+sdf.format(new Date()));
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated"); bb.writeLog("***** resource model initiated");
//
// try{ try{
// isMainIp = isMainIp();//计算是否为集群中的主节点 isMainIp = isMainIp();//计算是否为集群中的主节点
// isSampleMode = isSampleMode();//判断是否为单机 isSampleMode = isSampleMode();//判断是否为单机
// bb.writeLog("isMainIp:"+isMainIp); bb.writeLog("isMainIp:"+isMainIp);
// bb.writeLog("isSampleMode:"+isSampleMode); bb.writeLog("isSampleMode:"+isSampleMode);
//
// if(isSampleMode){ if(isSampleMode){
// initData(); initData();
// }else{ }else{
// if (isMainIp) { if (isMainIp) {
// //则调用service方法加载设置重新获取workflowid相关超时设置OvertimeEntity //则调用service方法加载设置重新获取workflowid相关超时设置OvertimeEntity
// initData(); initData();
// } }
// } }
//
//
//
// //DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("portal-producer-group_nj"); //DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("portal-producer-group_nj");
// //DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("cbec-consumer-group_nj_133"); //DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("cbec-consumer-group_nj_133");
// //consumer.setNamesrvAddr("172.16.84.183:9001;172.16.84.187:9001"); //consumer.setNamesrvAddr("172.16.84.183:9001;172.16.84.187:9001");
//// consumer.subscribe("dataSync_topic_nj", "BPM"); // consumer.subscribe("dataSync_topic_nj", "BPM");
//// consumer.setInstanceName("dataSync_topic_nj"); // consumer.setInstanceName("dataSync_topic_nj");
//
//// String consumerGroup = PropBean.getUfPropValue("consumerGroup"); // String consumerGroup = PropBean.getUfPropValue("consumerGroup");
//// bb.writeLog("consumerGroup:"+consumerGroup); // bb.writeLog("consumerGroup:"+consumerGroup);
//// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);
////
//// String namesrvAddr = PropBean.getUfPropValue("namesrvAddr");
//// bb.writeLog("namesrvAddr:"+namesrvAddr);
//// consumer.setNamesrvAddr(namesrvAddr);
////
//// String instanceName = PropBean.getUfPropValue("instanceName");
//// bb.writeLog("instanceName:"+instanceName);
//// consumer.setInstanceName(instanceName);
//// consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
//// //consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
////
//// String topic = PropBean.getUfPropValue("topic");
//// String subExpression = PropBean.getUfPropValue("subExpression");
////
//// bb.writeLog("topic:"+topic);
//// bb.writeLog("subExpression:"+subExpression);
////
//// consumer.subscribe(topic,subExpression);
//// consumer.setConsumeThreadMin(1);
//// consumer.setConsumeThreadMax(1);
//// consumer.setConsumeMessageBatchMaxSize(1);
//// consumer.setMessageModel(MessageModel.BROADCASTING);
//// consumer.registerMessageListener(new RocketMsgListener());
//// consumer.registerMessageListener(new RocketMsgOrderListener());
//
// /**
// * Consumer对象在使用之前必须要调用start初始化初始化一次即可<br>
// */
//// consumer.start();
//// bb.writeLog("Consumer Started.");
//
// }catch (Exception e){
// bb.writeLog("Consumer resource model initiated--MQClientException:"+e);
// }
//
// bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated");
// bb.writeLog("***** resource model initiated");
// }
//
//
// public void initData(){
// BaseBean bb = new BaseBean();
// PropBean propBean = new PropBean();
//
// try {
// String consumerGroup = propBean.getUfPropValueStatic("consumerGroup");
// bb.writeLog("consumerGroup:" + consumerGroup);
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup); // DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);
// String namesrvAddr = propBean.getUfPropValueStatic("namesrvAddr"); //
// bb.writeLog("namesrvAddr:" + namesrvAddr); // String namesrvAddr = PropBean.getUfPropValue("namesrvAddr");
// bb.writeLog("namesrvAddr:"+namesrvAddr);
// consumer.setNamesrvAddr(namesrvAddr); // consumer.setNamesrvAddr(namesrvAddr);
// String instanceName = propBean.getUfPropValueStatic("instanceName"); //
// bb.writeLog("instanceName:" + instanceName); // String instanceName = PropBean.getUfPropValue("instanceName");
// bb.writeLog("instanceName:"+instanceName);
// consumer.setInstanceName(instanceName); // consumer.setInstanceName(instanceName);
// consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); // consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// String topic = propBean.getUfPropValueStatic("topic"); // //consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
// String subExpression = propBean.getUfPropValueStatic("subExpression"); //
// bb.writeLog("topic:" + topic); // String topic = PropBean.getUfPropValue("topic");
// bb.writeLog("subExpression:" + subExpression); // String subExpression = PropBean.getUfPropValue("subExpression");
// String authid = propBean.getUfPropValueStatic("authid"); //
// String authpwd = propBean.getUfPropValueStatic("authpwd"); // bb.writeLog("topic:"+topic);
// String clustername = propBean.getUfPropValueStatic("clustername"); // bb.writeLog("subExpression:"+subExpression);
// String tenantid = propBean.getUfPropValueStatic("tenantid");
// bb.writeLog("authid:" + authid);
// bb.writeLog("authpwd:" + authpwd);
// bb.writeLog("clustername:" + clustername);
// bb.writeLog("tenantid:" + tenantid);
// consumer.subscribe(topic, subExpression);
// consumer.setAuthID(authid);
// consumer.setAuthPWD(authpwd);
// consumer.setClusterName(clustername);
// consumer.setTenantID(tenantid);
// //
// consumer.subscribe(topic,subExpression);
// consumer.setConsumeThreadMin(1); // consumer.setConsumeThreadMin(1);
// consumer.setConsumeThreadMax(1); // consumer.setConsumeThreadMax(1);
// consumer.setConsumeMessageBatchMaxSize(1); // consumer.setConsumeMessageBatchMaxSize(1);
//
// consumer.setMessageModel(MessageModel.BROADCASTING); // consumer.setMessageModel(MessageModel.BROADCASTING);
// consumer.registerMessageListener(new RocketMsgListener()); // consumer.registerMessageListener(new RocketMsgListener());
// consumer.registerMessageListener(new RocketMsgOrderListener());
/**
* Consumer使start<br>
*/
// consumer.start(); // consumer.start();
// bb.writeLog("Consumer66Started."); // bb.writeLog("Consumer Started.");
// } catch (MQClientException var13) {
// bb.writeLog("Consumer resource model initiated--MQClientException:" + var13); }catch (Exception e){
// } bb.writeLog("Consumer resource model initiated--MQClientException:"+e);
// } }
//
// /** bb.writeLog("***** resource model initiated");
// * 实现 HttpServlet 的 doGet 方法,不作任何操作 bb.writeLog("***** resource model initiated");
// * @param request bb.writeLog("***** resource model initiated");
// * @param response }
// * @throws ServletException
// * @throws IOException
// */ public void initData(){
// @Override BaseBean bb = new BaseBean();
// public void doGet(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException { PropBean propBean = new PropBean();
// doPost(request, response);
// } try {
// String consumerGroup = propBean.getUfPropValueStatic("consumerGroup");
// /** bb.writeLog("consumerGroup:" + consumerGroup);
// * 实现 HttpServlet 的 doPost 方法,不作任何操作 DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);
// * @param request String namesrvAddr = propBean.getUfPropValueStatic("namesrvAddr");
// * @param response bb.writeLog("namesrvAddr:" + namesrvAddr);
// * @throws ServletException consumer.setNamesrvAddr(namesrvAddr);
// * @throws IOException String instanceName = propBean.getUfPropValueStatic("instanceName");
// */ bb.writeLog("instanceName:" + instanceName);
// consumer.setInstanceName(instanceName);
// @Override consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// public void doPost(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException { String topic = propBean.getUfPropValueStatic("topic");
// } String subExpression = propBean.getUfPropValueStatic("subExpression");
// bb.writeLog("topic:" + topic);
// @Override bb.writeLog("subExpression:" + subExpression);
// public void destroy() String authid = propBean.getUfPropValueStatic("authid");
// { String authpwd = propBean.getUfPropValueStatic("authpwd");
// // 什么也不做 String clustername = propBean.getUfPropValueStatic("clustername");
// } String tenantid = propBean.getUfPropValueStatic("tenantid");
// bb.writeLog("authid:" + authid);
// bb.writeLog("authpwd:" + authpwd);
// private boolean isMainIp() { bb.writeLog("clustername:" + clustername);
// new BaseBean().writeLog("超时判断主节点"); bb.writeLog("tenantid:" + tenantid);
// String mainControlIp = ""; consumer.subscribe(topic, subExpression);
// ArrayList<String> hostIps = new InitServer().getRealIp(); consumer.setAuthID(authid);
// Prop prop = Prop.getInstance(); consumer.setAuthPWD(authpwd);
// mainControlIp = prop.getPropValue(GCONST.getConfigFile(), "MainControlIP"); consumer.setClusterName(clustername);
// if (hostIps == null || hostIps.size() == 0) { consumer.setTenantID(tenantid);
// new BaseBean().writeLog("System Init Error:Cannot get local Ip address,This may cause scripts or Timed task not run! ");
// } else { consumer.setConsumeThreadMin(1);
// new BaseBean().writeLog("System Init Message:mainControlIp=" + mainControlIp + " localIp:" + hostIps.toString()); consumer.setConsumeThreadMax(1);
// } consumer.setConsumeMessageBatchMaxSize(1);
// if ((!"".equals(mainControlIp) && hostIps.contains(mainControlIp)) || "".equals(mainControlIp)) {
// return true; consumer.setMessageModel(MessageModel.BROADCASTING);
// } consumer.registerMessageListener(new RocketMsgListener());
// return false; consumer.start();
// } bb.writeLog("Consumer66Started.");
// } catch (MQClientException var13) {
// public boolean isSampleMode() { bb.writeLog("Consumer resource model initiated--MQClientException:" + var13);
// BaseBean base = new BaseBean(); }
// boolean redis_flag = "1".equals(base.getPropValue("weaver_new_session", "status")); }
// boolean mainIp_flag = !"".equals(base.getPropValue("weaver", "MainControlIP"));
// base.writeLog("超时判断是否为集群环境redis_flag======"+redis_flag); /**
// base.writeLog("超时判断是否为集群环境mainIp_flag======="+mainIp_flag); * HttpServlet doGet
// return !(redis_flag && mainIp_flag); * @param request
// } * @param response
// * @throws ServletException
// * @throws IOException
// */
//} @Override
public void doGet(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {
doPost(request, response);
}
/**
* HttpServlet doPost
* @param request
* @param response
* @throws ServletException
* @throws IOException
*/
@Override
public void doPost(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {
}
@Override
public void destroy()
{
// 什么也不做
}
private boolean isMainIp() {
new BaseBean().writeLog("超时判断主节点");
String mainControlIp = "";
ArrayList<String> hostIps = new InitServer().getRealIp();
Prop prop = Prop.getInstance();
mainControlIp = prop.getPropValue(GCONST.getConfigFile(), "MainControlIP");
if (hostIps == null || hostIps.size() == 0) {
new BaseBean().writeLog("System Init Error:Cannot get local Ip address,This may cause scripts or Timed task not run! ");
} else {
new BaseBean().writeLog("System Init Message:mainControlIp=" + mainControlIp + " localIp:" + hostIps.toString());
}
if ((!"".equals(mainControlIp) && hostIps.contains(mainControlIp)) || "".equals(mainControlIp)) {
return true;
}
return false;
}
public boolean isSampleMode() {
BaseBean base = new BaseBean();
boolean redis_flag = "1".equals(base.getPropValue("weaver_new_session", "status"));
boolean mainIp_flag = !"".equals(base.getPropValue("weaver", "MainControlIP"));
base.writeLog("超时判断是否为集群环境redis_flag======"+redis_flag);
base.writeLog("超时判断是否为集群环境mainIp_flag======="+mainIp_flag);
return !(redis_flag && mainIp_flag);
}
}

@ -14,14 +14,18 @@ import weaver.hrm.resource.ResourceComInfo;
import weaver.interfaces.dito.comInfo.PropBean; import weaver.interfaces.dito.comInfo.PropBean;
import java.text.SimpleDateFormat; import java.text.SimpleDateFormat;
import java.util.Date; import java.util.*;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
public class RocketmqUtil { public class RocketmqUtil {
private List<ConsumerBase> consumerBases;
{
consumerBases = new ArrayList<>();
consumerBases.add(new SignatureConsumer());
}
private Lock lock = new ReentrantLock(); private Lock lock = new ReentrantLock();
public int updateOrgData(String data) public int updateOrgData(String data)
{ {
@ -58,6 +62,11 @@ public class RocketmqUtil {
}else if("staff".equals(tableName)){ }else if("staff".equals(tableName)){
updateStaffData(jsonArray,tableName); updateStaffData(jsonArray,tableName);
} }
for (ConsumerBase consumerBase :consumerBases){
if (consumerBase.support(tableName)){
consumerBase.consumer(jsonArray,tableName);
}
}
} }
} }
} }

@ -0,0 +1,32 @@
package weaver.interfaces.dito.mq;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import weaver.general.BaseBean;
import weaver.interfaces.dito.job.ESiginsCronJob;
import java.util.Arrays;
public class SignatureConsumer implements ConsumerBase{
@Override
public boolean support(String tableName) {
return "ele_signature".equals(tableName);
}
@Override
public void consumer(JSONArray jsonArray, String tableName) {
BaseBean bb = new BaseBean();
bb.writeLog("********** consumer SignatureConsumer start **********");
for (int i = 0; i < jsonArray.size(); i++){
JSONObject jsonObject = jsonArray.getJSONObject(i);
String sysUserCode = jsonObject.getString("sysUserCode");
ESiginsCronJob eSiginsCronJob = new ESiginsCronJob();
eSiginsCronJob.setUsercode(sysUserCode);
eSiginsCronJob.execute();
}
bb.writeLog("********** consumer SignatureConsumer end **********");
}
}
Loading…
Cancel
Save