Prechádzať zdrojové kódy

推送数据 加锁

miaolijing 3 rokov pred
rodič
commit
8b3e51c66b

+ 0 - 16
src/main/java/com/persagy/apm/diagnose/config/RabbitConfig.java

@@ -31,13 +31,10 @@ public class RabbitConfig {
     private String password;
 
     //交换机
-    public static final String EXCHANGE_MAINTENANCE = "exchange_maintenance";
     public static final String EXCHANGE_INDICATOR = "exchange_indicator";
     //队列
-    public static final String QUEUE_MAINTENANCE = "QUEUE_maintenance";
     public static final String QUEUE_INDICATOR = "QUEUE_indicator";
     //路由键
-    public static final String ROUTINGKEY_MAINTENANCE = "routingKey_maintenance";
     public static final String ROUTINGKEY_INDICATOR = "routingKey_indicator";
 
     @Bean
@@ -84,17 +81,4 @@ public class RabbitConfig {
         return BindingBuilder.bind(queueIndicator()).to(defaultExchange()).with(RabbitConfig.ROUTINGKEY_INDICATOR);
     }
 
-    /**
-     * 获取队列B
-     * @return
-     */
-    @Bean
-    public Queue queueMaintenance() {
-        return new Queue(QUEUE_MAINTENANCE, true); //队列持久
-    }
-
-    @Bean
-    public Binding bindingMaintenance() {
-        return BindingBuilder.bind(queueMaintenance()).to(defaultExchange()).with(RabbitConfig.ROUTINGKEY_MAINTENANCE);
-    }
 }

+ 3 - 0
src/main/java/com/persagy/apm/diagnose/indicatorrecord/service/impl/MonitorIndicatorRecordServiceImpl.java

@@ -392,6 +392,9 @@ public class MonitorIndicatorRecordServiceImpl implements IMonitorIndicatorRecor
             }
             JSONArray sendArray = CollectDataUtil.batchBuildSendJsonParam(sendTimeKeyAndDataList,alarmItemCode);
             msgProducer.sendIndicatorMsg(sendArray);
+            //加锁
+            long time1 = System.currentTimeMillis() + (20 * 1000);
+            lockUtil.lock(projectDTO.getProjectId()+objIdAndAlarmItemCode+ "_sendData", String.valueOf(time1));
 //          String sentValue = CollectDataUtil.batchBuildSendParam(sendTimeKeyAndDataList,alarmItemCode);
             //AlarmWebSocketServer.sendMsgToClients(projectDTO.getProjectId(), sentValue);
             log.error("指标发送报表服务数据:" + projectDTO.getProjectId()+";"+ sendArray);

+ 4 - 1
src/main/java/com/persagy/apm/diagnose/maintenance/service/impl/ProjectDataRecordServiceImpl.java

@@ -318,7 +318,10 @@ public class ProjectDataRecordServiceImpl implements IProjectDataRecordService {
                 endTime = com.persagy.apm.diagnose.utils.DateUtils.str2Date(dateListEntry.getKey(), com.persagy.apm.diagnose.utils.DateUtils.SDFSECOND);
             }
             JSONArray sendArray = CollectDataUtil.batchBuildSendJsonParam(sendTimeKeyAndDataList,alarmItemCode);
-            msgProducer.sendMaintenanceMsg(sendArray);
+            msgProducer.sendIndicatorMsg(sendArray);
+            //加锁
+            long time1 = System.currentTimeMillis() + (20 * 1000);
+            lockUtil.lock(projectId +objIdAndAlarmItemCode+ "_sendMaintenance", String.valueOf(time1));
 //            String sentValue = CollectDataUtil.batchBuildSendParam(sendTimeKeyAndDataList,alarmItemCode);
          //   AlarmWebSocketServer.sendMsgToClients(projectId, sentValue);
             log.info("设备维保发送数据服务数据:" + projectId+";"+ sendArray);

+ 0 - 6
src/main/java/com/persagy/apm/diagnose/service/MsgProducer.java

@@ -36,12 +36,6 @@ public class MsgProducer implements RabbitTemplate.ConfirmCallback {
         rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_INDICATOR, RabbitConfig.ROUTINGKEY_INDICATOR, contentArray, correlationId);
     }
 
-    public void sendMaintenanceMsg(JSONArray contentArray) {
-        CorrelationData correlationId = new CorrelationData(UUID.randomUUID().toString());
-        //把消息放入routingKey_maintenance对应的队列当中去,对应的是队列QUEUE_maintenance
-        rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_MAINTENANCE, RabbitConfig.ROUTINGKEY_MAINTENANCE,contentArray, correlationId);
-    }
-
     /**
      * 回调
      */