| 
					
				 | 
			
			
				@@ -2,6 +2,7 @@ package com.genersoft.iot.vmp.service.redisMsg; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				  
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import com.alibaba.fastjson2.JSON; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import com.alibaba.fastjson2.JSONObject; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+import com.genersoft.iot.vmp.gb28181.bean.GbStream; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import com.genersoft.iot.vmp.service.IGbStreamService; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import com.genersoft.iot.vmp.service.IMediaServerService; 
			 | 
		
	
	
		
			
				| 
					
				 | 
			
			
				@@ -19,6 +20,7 @@ import org.springframework.stereotype.Component; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import javax.annotation.Resource; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import java.util.ArrayList; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import java.util.List; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+import java.util.Map; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 import java.util.concurrent.ConcurrentLinkedQueue; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				  
			 | 
		
	
		
			
				 | 
				 | 
			
			
				 /** 
			 | 
		
	
	
		
			
				| 
					
				 | 
			
			
				@@ -57,7 +59,8 @@ public class RedisPushStreamStatusListMsgListener implements MessageListener { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                     try { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         List<StreamPushItem> streamPushItems = JSON.parseArray(new String(msg.getBody()), StreamPushItem.class); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         //查询全部的app+stream 用于判断是添加还是修改 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				-                        List<String> allAppAndStream = streamPushService.getAllAppAndStream(); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                        Map<String, StreamPushItem> allAppAndStream = streamPushService.getAllAppAndStreamMap(); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                        Map<String, GbStream> allGBId = gbStreamService.getAllGBId(); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				  
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         /** 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                          * 用于存储更具APP+Stream过滤后的数据,可以直接存入stream_push表与gb_stream表 
			 | 
		
	
	
		
			
				| 
					
				 | 
			
			
				@@ -67,9 +70,15 @@ public class RedisPushStreamStatusListMsgListener implements MessageListener { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         for (StreamPushItem streamPushItem : streamPushItems) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             String app = streamPushItem.getApp(); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             String stream = streamPushItem.getStream(); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				-                            boolean contains = allAppAndStream.contains(app + stream); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                            boolean contains = allAppAndStream.containsKey(app + stream); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             //不存在就添加 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             if (!contains) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                if (allGBId.containsKey(streamPushItem.getGbId())) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                    GbStream gbStream = allGBId.get(streamPushItem.getGbId()); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                    logger.warn("[REDIS消息-推流设备列表更新] 国标编号重复: {}, 已分配给{}/{}", 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                            streamPushItem.getGbId(), gbStream.getApp(), gbStream.getStream()); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                    continue; 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItem.setStreamType("push"); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItem.setCreateTime(DateUtil.getNow()); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItem.setMediaServerId(mediaServerService.getDefaultMediaServer().getId()); 
			 | 
		
	
	
		
			
				| 
					
				 | 
			
			
				@@ -77,25 +86,25 @@ public class RedisPushStreamStatusListMsgListener implements MessageListener { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItem.setOriginTypeStr("rtsp_push"); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItem.setTotalReaderCount("0"); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItemForSave.add(streamPushItem); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                                allGBId.put(streamPushItem.getGbId(), streamPushItem); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             } else { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 //存在就只修改 name和gbId 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                                 streamPushItemForUpdate.add(streamPushItem); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				-                        if (streamPushItemForSave.size() > 0) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				- 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                        if (!streamPushItemForSave.isEmpty()) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             logger.info("添加{}条",streamPushItemForSave.size()); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             logger.info(JSONObject.toJSONString(streamPushItemForSave)); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             streamPushService.batchAdd(streamPushItemForSave); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				  
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				-                        if(streamPushItemForUpdate.size()>0){ 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                        if(!streamPushItemForUpdate.isEmpty()){ 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             logger.info("修改{}条",streamPushItemForUpdate.size()); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             logger.info(JSONObject.toJSONString(streamPushItemForUpdate)); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                             gbStreamService.updateGbIdOrName(streamPushItemForUpdate); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                     }catch (Exception e) { 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				-                        logger.warn("[REDIS消息-推流设备列表更新] 发现未处理的异常, \r\n{}", JSON.toJSONString(message)); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				+                        logger.warn("[REDIS消息-推流设备列表更新] 发现未处理的异常, \r\n{}", new String(message.getBody())); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                         logger.error("[REDIS消息-推流设备列表更新] 异常内容: ", e); 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                     } 
			 | 
		
	
		
			
				 | 
				 | 
			
			
				                 } 
			 |