StreamPushServiceImpl.java 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636
  1. package com.genersoft.iot.vmp.service.impl;
  2. import com.alibaba.fastjson2.JSONObject;
  3. import com.baomidou.dynamic.datasource.annotation.DS;
  4. import com.genersoft.iot.vmp.common.StreamInfo;
  5. import com.genersoft.iot.vmp.conf.MediaConfig;
  6. import com.genersoft.iot.vmp.conf.UserSetting;
  7. import com.genersoft.iot.vmp.gb28181.bean.GbStream;
  8. import com.genersoft.iot.vmp.gb28181.bean.ParentPlatform;
  9. import com.genersoft.iot.vmp.gb28181.bean.PlatformCatalog;
  10. import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
  11. import com.genersoft.iot.vmp.gb28181.event.subscribe.catalog.CatalogEvent;
  12. import com.genersoft.iot.vmp.media.bean.MediaInfo;
  13. import com.genersoft.iot.vmp.media.event.MediaArrivalEvent;
  14. import com.genersoft.iot.vmp.media.event.MediaDepartureEvent;
  15. import com.genersoft.iot.vmp.media.service.IMediaServerService;
  16. import com.genersoft.iot.vmp.media.zlm.dto.MediaServer;
  17. import com.genersoft.iot.vmp.media.zlm.dto.StreamAuthorityInfo;
  18. import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem;
  19. import com.genersoft.iot.vmp.media.zlm.dto.hook.OnStreamChangedHookParam;
  20. import com.genersoft.iot.vmp.media.zlm.dto.hook.OriginType;
  21. import com.genersoft.iot.vmp.service.IGbStreamService;
  22. import com.genersoft.iot.vmp.service.IStreamPushService;
  23. import com.genersoft.iot.vmp.service.bean.StreamPushItemFromRedis;
  24. import com.genersoft.iot.vmp.storager.IRedisCatchStorage;
  25. import com.genersoft.iot.vmp.storager.dao.*;
  26. import com.genersoft.iot.vmp.utils.DateUtil;
  27. import com.genersoft.iot.vmp.vmanager.bean.ResourceBaseInfo;
  28. import com.github.pagehelper.PageHelper;
  29. import com.github.pagehelper.PageInfo;
  30. import org.slf4j.Logger;
  31. import org.slf4j.LoggerFactory;
  32. import org.springframework.beans.factory.annotation.Autowired;
  33. import org.springframework.context.event.EventListener;
  34. import org.springframework.jdbc.datasource.DataSourceTransactionManager;
  35. import org.springframework.scheduling.annotation.Async;
  36. import org.springframework.stereotype.Service;
  37. import org.springframework.transaction.TransactionDefinition;
  38. import org.springframework.transaction.TransactionStatus;
  39. import org.springframework.util.ObjectUtils;
  40. import java.util.*;
  41. import java.util.stream.Collectors;
  42. @Service
  43. @DS("master")
  44. public class StreamPushServiceImpl implements IStreamPushService {
  45. private final static Logger logger = LoggerFactory.getLogger(StreamPushServiceImpl.class);
  46. @Autowired
  47. private GbStreamMapper gbStreamMapper;
  48. @Autowired
  49. private StreamPushMapper streamPushMapper;
  50. @Autowired
  51. private StreamProxyMapper streamProxyMapper;
  52. @Autowired
  53. private ParentPlatformMapper parentPlatformMapper;
  54. @Autowired
  55. private PlatformCatalogMapper platformCatalogMapper;
  56. @Autowired
  57. private PlatformGbStreamMapper platformGbStreamMapper;
  58. @Autowired
  59. private IGbStreamService gbStreamService;
  60. @Autowired
  61. private EventPublisher eventPublisher;
  62. @Autowired
  63. private IRedisCatchStorage redisCatchStorage;
  64. @Autowired
  65. private UserSetting userSetting;
  66. @Autowired
  67. private IMediaServerService mediaServerService;
  68. @Autowired
  69. DataSourceTransactionManager dataSourceTransactionManager;
  70. @Autowired
  71. TransactionDefinition transactionDefinition;
  72. @Autowired
  73. private MediaConfig mediaConfig;
  74. /**
  75. * 流到来的处理
  76. */
  77. @Async("taskExecutor")
  78. @EventListener
  79. public void onApplicationEvent(MediaArrivalEvent event) {
  80. MediaInfo mediaInfo = event.getMediaInfo();
  81. if (mediaInfo == null) {
  82. return;
  83. }
  84. if (mediaInfo.getOriginType() != OriginType.RTMP_PUSH.ordinal()
  85. && mediaInfo.getOriginType() != OriginType.RTSP_PUSH.ordinal()
  86. && mediaInfo.getOriginType() != OriginType.RTC_PUSH.ordinal()) {
  87. return;
  88. }
  89. StreamAuthorityInfo streamAuthorityInfo = redisCatchStorage.getStreamAuthorityInfo(event.getApp(), event.getStream());
  90. if (streamAuthorityInfo == null) {
  91. streamAuthorityInfo = StreamAuthorityInfo.getInstanceByHook(event);
  92. } else {
  93. streamAuthorityInfo.setOriginType(mediaInfo.getOriginType());
  94. }
  95. redisCatchStorage.updateStreamAuthorityInfo(event.getApp(), event.getStream(), streamAuthorityInfo);
  96. StreamPushItem transform = StreamPushItem.getInstance(event, userSetting.getServerId());
  97. transform.setPushIng(true);
  98. transform.setUpdateTime(DateUtil.getNow());
  99. transform.setPushTime(DateUtil.getNow());
  100. transform.setSelf(true);
  101. StreamPushItem pushInDb = getPush(event.getApp(), event.getStream());
  102. if (pushInDb == null) {
  103. transform.setCreateTime(DateUtil.getNow());
  104. streamPushMapper.add(transform);
  105. }else {
  106. streamPushMapper.update(transform);
  107. gbStreamMapper.updateMediaServer(event.getApp(), event.getStream(), event.getMediaServer().getId());
  108. }
  109. // TODO 相关的事件自行管理,不需要写入ZLMMediaListManager
  110. // ChannelOnlineEvent channelOnlineEventLister = getChannelOnlineEventLister(transform.getApp(), transform.getStream());
  111. // if ( channelOnlineEventLister != null) {
  112. // try {
  113. // channelOnlineEventLister.run(transform.getApp(), transform.getStream(), transform.getServerId());;
  114. // } catch (ParseException e) {
  115. // logger.error("addPush: ", e);
  116. // }
  117. // removedChannelOnlineEventLister(transform.getApp(), transform.getStream());
  118. // }
  119. // 冗余数据,自己系统中自用
  120. redisCatchStorage.addPushListItem(event.getApp(), event.getStream(), event);
  121. // 发送流变化redis消息
  122. JSONObject jsonObject = new JSONObject();
  123. jsonObject.put("serverId", userSetting.getServerId());
  124. jsonObject.put("app", event.getApp());
  125. jsonObject.put("stream", event.getStream());
  126. jsonObject.put("register", true);
  127. jsonObject.put("mediaServerId", event.getMediaServer().getId());
  128. redisCatchStorage.sendStreamChangeMsg(OriginType.values()[event.getMediaInfo().getOriginType()].getType(), jsonObject);
  129. }
  130. /**
  131. * 流离开的处理
  132. */
  133. @Async("taskExecutor")
  134. @EventListener
  135. public void onApplicationEvent(MediaDepartureEvent event) {
  136. // 兼容流注销时类型从redis记录获取
  137. OnStreamChangedHookParam onStreamChangedHookParam = redisCatchStorage.getStreamInfo(
  138. event.getApp(), event.getStream(), event.getMediaServer().getId());
  139. if (onStreamChangedHookParam != null) {
  140. String type = OriginType.values()[onStreamChangedHookParam.getOriginType()].getType();
  141. redisCatchStorage.removeStream(event.getMediaServer().getId(), type, event.getApp(), event.getStream());
  142. if ("PUSH".equalsIgnoreCase(type)) {
  143. // 冗余数据,自己系统中自用
  144. redisCatchStorage.removePushListItem(event.getApp(), event.getStream(), event.getMediaServer().getId());
  145. }
  146. if (type != null) {
  147. // 发送流变化redis消息
  148. JSONObject jsonObject = new JSONObject();
  149. jsonObject.put("serverId", userSetting.getServerId());
  150. jsonObject.put("app", event.getApp());
  151. jsonObject.put("stream", event.getStream());
  152. jsonObject.put("register", false);
  153. jsonObject.put("mediaServerId", event.getMediaServer().getId());
  154. redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
  155. }
  156. }
  157. GbStream gbStream = gbStreamMapper.selectOne(event.getApp(), event.getStream());
  158. if (gbStream != null) {
  159. if (userSetting.isUsePushingAsStatus()) {
  160. streamPushMapper.updatePushStatus(event.getApp(), event.getStream(), false);
  161. eventPublisher.catalogEventPublishForStream(null, gbStream, CatalogEvent.OFF);
  162. }
  163. }else {
  164. streamPushMapper.del(event.getApp(), event.getStream());
  165. }
  166. }
  167. private List<StreamPushItem> handleJSON(List<StreamInfo> streamInfoList) {
  168. if (streamInfoList == null || streamInfoList.isEmpty()) {
  169. return null;
  170. }
  171. Map<String, StreamPushItem> result = new HashMap<>();
  172. for (StreamInfo streamInfo : streamInfoList) {
  173. // 不保存国标推理以及拉流代理的流
  174. if (streamInfo.getOriginType() == OriginType.RTSP_PUSH.ordinal()
  175. || streamInfo.getOriginType() == OriginType.RTMP_PUSH.ordinal()
  176. || streamInfo.getOriginType() == OriginType.RTC_PUSH.ordinal() ) {
  177. String key = streamInfo.getApp() + "_" + streamInfo.getStream();
  178. StreamPushItem streamPushItem = result.get(key);
  179. if (streamPushItem == null) {
  180. streamPushItem = streamPushItem.getInstance(streamInfo);
  181. result.put(key, streamPushItem);
  182. }
  183. }
  184. }
  185. return new ArrayList<>(result.values());
  186. }
  187. @Override
  188. public StreamPushItem transform(OnStreamChangedHookParam item) {
  189. StreamPushItem streamPushItem = new StreamPushItem();
  190. streamPushItem.setApp(item.getApp());
  191. streamPushItem.setMediaServerId(item.getMediaServerId());
  192. streamPushItem.setStream(item.getStream());
  193. streamPushItem.setAliveSecond(item.getAliveSecond());
  194. streamPushItem.setOriginSock(item.getOriginSock());
  195. streamPushItem.setTotalReaderCount(item.getTotalReaderCount() + "");
  196. streamPushItem.setOriginType(item.getOriginType());
  197. streamPushItem.setOriginTypeStr(item.getOriginTypeStr());
  198. streamPushItem.setOriginUrl(item.getOriginUrl());
  199. streamPushItem.setCreateTime(DateUtil.getNow());
  200. streamPushItem.setAliveSecond(item.getAliveSecond());
  201. streamPushItem.setStatus(true);
  202. streamPushItem.setStreamType("push");
  203. streamPushItem.setVhost(item.getVhost());
  204. streamPushItem.setServerId(item.getSeverId());
  205. return streamPushItem;
  206. }
  207. @Override
  208. public PageInfo<StreamPushItem> getPushList(Integer page, Integer count, String query, Boolean pushing, String mediaServerId) {
  209. PageHelper.startPage(page, count);
  210. List<StreamPushItem> all = streamPushMapper.selectAllForList(query, pushing, mediaServerId);
  211. return new PageInfo<>(all);
  212. }
  213. @Override
  214. public List<StreamPushItem> getPushList(String mediaServerId) {
  215. return streamPushMapper.selectAllByMediaServerIdWithOutGbID(mediaServerId);
  216. }
  217. @Override
  218. public boolean saveToGB(GbStream stream) {
  219. stream.setStreamType("push");
  220. stream.setStatus(true);
  221. stream.setCreateTime(DateUtil.getNow());
  222. stream.setStreamType("push");
  223. stream.setMediaServerId(mediaConfig.getId());
  224. int add = gbStreamMapper.add(stream);
  225. return add > 0;
  226. }
  227. @Override
  228. public boolean removeFromGB(GbStream stream) {
  229. // 判断是否需要发送事件
  230. gbStreamService.sendCatalogMsg(stream, CatalogEvent.DEL);
  231. platformGbStreamMapper.delByAppAndStream(stream.getApp(), stream.getStream());
  232. int del = gbStreamMapper.del(stream.getApp(), stream.getStream());
  233. MediaServer mediaInfo = mediaServerService.getOne(stream.getMediaServerId());
  234. List<StreamInfo> mediaList = mediaServerService.getMediaList(mediaInfo, stream.getApp(), stream.getStream(), null);
  235. if (mediaList != null && mediaList.isEmpty()) {
  236. streamPushMapper.del(stream.getApp(), stream.getStream());
  237. }
  238. return del > 0;
  239. }
  240. @Override
  241. public StreamPushItem getPush(String app, String streamId) {
  242. return streamPushMapper.selectOne(app, streamId);
  243. }
  244. @Override
  245. public boolean stop(String app, String streamId) {
  246. logger.info("[推流 ] 停止流: {}/{}", app, streamId);
  247. StreamPushItem streamPushItem = streamPushMapper.selectOne(app, streamId);
  248. if (streamPushItem != null) {
  249. gbStreamService.sendCatalogMsg(streamPushItem, CatalogEvent.DEL);
  250. }
  251. platformGbStreamMapper.delByAppAndStream(app, streamId);
  252. gbStreamMapper.del(app, streamId);
  253. int delStream = streamPushMapper.del(app, streamId);
  254. if (delStream > 0) {
  255. MediaServer mediaServerItem = mediaServerService.getOne(streamPushItem.getMediaServerId());
  256. mediaServerService.closeStreams(mediaServerItem,app, streamId);
  257. }
  258. return true;
  259. }
  260. @Override
  261. public void zlmServerOnline(String mediaServerId) {
  262. // 同步zlm推流信息
  263. MediaServer mediaServerItem = mediaServerService.getOne(mediaServerId);
  264. if (mediaServerItem == null) {
  265. return;
  266. }
  267. // 数据库记录
  268. List<StreamPushItem> pushList = getPushList(mediaServerId);
  269. Map<String, StreamPushItem> pushItemMap = new HashMap<>();
  270. // redis记录
  271. List<OnStreamChangedHookParam> onStreamChangedHookParams = redisCatchStorage.getStreams(mediaServerId, "PUSH");
  272. Map<String, OnStreamChangedHookParam> streamInfoPushItemMap = new HashMap<>();
  273. if (pushList.size() > 0) {
  274. for (StreamPushItem streamPushItem : pushList) {
  275. if (ObjectUtils.isEmpty(streamPushItem.getGbId())) {
  276. pushItemMap.put(streamPushItem.getApp() + streamPushItem.getStream(), streamPushItem);
  277. }
  278. }
  279. }
  280. if (onStreamChangedHookParams.size() > 0) {
  281. for (OnStreamChangedHookParam onStreamChangedHookParam : onStreamChangedHookParams) {
  282. streamInfoPushItemMap.put(onStreamChangedHookParam.getApp() + onStreamChangedHookParam.getStream(), onStreamChangedHookParam);
  283. }
  284. }
  285. // 获取所有推流鉴权信息,清理过期的
  286. List<StreamAuthorityInfo> allStreamAuthorityInfo = redisCatchStorage.getAllStreamAuthorityInfo();
  287. Map<String, StreamAuthorityInfo> streamAuthorityInfoInfoMap = new HashMap<>();
  288. for (StreamAuthorityInfo streamAuthorityInfo : allStreamAuthorityInfo) {
  289. streamAuthorityInfoInfoMap.put(streamAuthorityInfo.getApp() + streamAuthorityInfo.getStream(), streamAuthorityInfo);
  290. }
  291. List<StreamInfo> mediaList = mediaServerService.getMediaList(mediaServerItem, null, null, null);
  292. if (mediaList == null) {
  293. return;
  294. }
  295. List<StreamPushItem> streamPushItems = handleJSON(mediaList);
  296. if (streamPushItems != null) {
  297. for (StreamPushItem streamPushItem : streamPushItems) {
  298. pushItemMap.remove(streamPushItem.getApp() + streamPushItem.getStream());
  299. streamInfoPushItemMap.remove(streamPushItem.getApp() + streamPushItem.getStream());
  300. streamAuthorityInfoInfoMap.remove(streamPushItem.getApp() + streamPushItem.getStream());
  301. }
  302. }
  303. List<StreamPushItem> offlinePushItems = new ArrayList<>(pushItemMap.values());
  304. if (offlinePushItems.size() > 0) {
  305. String type = "PUSH";
  306. int runLimit = 300;
  307. if (offlinePushItems.size() > runLimit) {
  308. for (int i = 0; i < offlinePushItems.size(); i += runLimit) {
  309. int toIndex = i + runLimit;
  310. if (i + runLimit > offlinePushItems.size()) {
  311. toIndex = offlinePushItems.size();
  312. }
  313. List<StreamPushItem> streamPushItemsSub = offlinePushItems.subList(i, toIndex);
  314. streamPushMapper.delAll(streamPushItemsSub);
  315. }
  316. }else {
  317. streamPushMapper.delAll(offlinePushItems);
  318. }
  319. }
  320. Collection<OnStreamChangedHookParam> offlineOnStreamChangedHookParamList = streamInfoPushItemMap.values();
  321. if (offlineOnStreamChangedHookParamList.size() > 0) {
  322. String type = "PUSH";
  323. for (OnStreamChangedHookParam offlineOnStreamChangedHookParam : offlineOnStreamChangedHookParamList) {
  324. JSONObject jsonObject = new JSONObject();
  325. jsonObject.put("serverId", userSetting.getServerId());
  326. jsonObject.put("app", offlineOnStreamChangedHookParam.getApp());
  327. jsonObject.put("stream", offlineOnStreamChangedHookParam.getStream());
  328. jsonObject.put("register", false);
  329. jsonObject.put("mediaServerId", mediaServerId);
  330. redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
  331. // 移除redis内流的信息
  332. redisCatchStorage.removeStream(mediaServerItem.getId(), "PUSH", offlineOnStreamChangedHookParam.getApp(), offlineOnStreamChangedHookParam.getStream());
  333. // 冗余数据,自己系统中自用
  334. redisCatchStorage.removePushListItem(offlineOnStreamChangedHookParam.getApp(), offlineOnStreamChangedHookParam.getStream(), mediaServerItem.getId());
  335. }
  336. }
  337. Collection<StreamAuthorityInfo> streamAuthorityInfos = streamAuthorityInfoInfoMap.values();
  338. if (streamAuthorityInfos.size() > 0) {
  339. for (StreamAuthorityInfo streamAuthorityInfo : streamAuthorityInfos) {
  340. // 移除redis内流的信息
  341. redisCatchStorage.removeStreamAuthorityInfo(streamAuthorityInfo.getApp(), streamAuthorityInfo.getStream());
  342. }
  343. }
  344. }
  345. @Override
  346. public void zlmServerOffline(String mediaServerId) {
  347. List<StreamPushItem> streamPushItems = streamPushMapper.selectAllByMediaServerIdWithOutGbID(mediaServerId);
  348. // 移除没有GBId的推流
  349. streamPushMapper.deleteWithoutGBId(mediaServerId);
  350. gbStreamMapper.deleteWithoutGBId("push", mediaServerId);
  351. // 其他的流设置未启用
  352. streamPushMapper.updateStatusByMediaServerId(mediaServerId, false);
  353. streamProxyMapper.updateStatusByMediaServerId(mediaServerId, false);
  354. // 发送流停止消息
  355. String type = "PUSH";
  356. // 发送redis消息
  357. List<OnStreamChangedHookParam> streamInfoList = redisCatchStorage.getStreams(mediaServerId, type);
  358. if (streamInfoList.size() > 0) {
  359. for (OnStreamChangedHookParam onStreamChangedHookParam : streamInfoList) {
  360. // 移除redis内流的信息
  361. redisCatchStorage.removeStream(mediaServerId, type, onStreamChangedHookParam.getApp(), onStreamChangedHookParam.getStream());
  362. JSONObject jsonObject = new JSONObject();
  363. jsonObject.put("serverId", userSetting.getServerId());
  364. jsonObject.put("app", onStreamChangedHookParam.getApp());
  365. jsonObject.put("stream", onStreamChangedHookParam.getStream());
  366. jsonObject.put("register", false);
  367. jsonObject.put("mediaServerId", mediaServerId);
  368. redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
  369. // 冗余数据,自己系统中自用
  370. redisCatchStorage.removePushListItem(onStreamChangedHookParam.getApp(), onStreamChangedHookParam.getStream(), mediaServerId);
  371. }
  372. }
  373. }
  374. @Override
  375. public void clean() {
  376. }
  377. @Override
  378. public boolean saveToRandomGB() {
  379. List<StreamPushItem> streamPushItems = streamPushMapper.selectAll();
  380. long gbId = 100001;
  381. for (StreamPushItem streamPushItem : streamPushItems) {
  382. streamPushItem.setStreamType("push");
  383. streamPushItem.setStatus(true);
  384. streamPushItem.setGbId("34020000004111" + gbId);
  385. streamPushItem.setCreateTime(DateUtil.getNow());
  386. gbId ++;
  387. }
  388. int limitCount = 30;
  389. if (streamPushItems.size() > limitCount) {
  390. for (int i = 0; i < streamPushItems.size(); i += limitCount) {
  391. int toIndex = i + limitCount;
  392. if (i + limitCount > streamPushItems.size()) {
  393. toIndex = streamPushItems.size();
  394. }
  395. gbStreamMapper.batchAdd(streamPushItems.subList(i, toIndex));
  396. }
  397. }else {
  398. gbStreamMapper.batchAdd(streamPushItems);
  399. }
  400. return true;
  401. }
  402. @Override
  403. public void batchAdd(List<StreamPushItem> streamPushItems) {
  404. streamPushMapper.addAll(streamPushItems);
  405. gbStreamMapper.batchAdd(streamPushItems);
  406. }
  407. @Override
  408. public void batchAddForUpload(List<StreamPushItem> streamPushItems, Map<String, List<String[]>> streamPushItemsForAll ) {
  409. // 存储数据到stream_push表
  410. streamPushMapper.addAll(streamPushItems);
  411. List<StreamPushItem> streamPushItemForGbStream = streamPushItems.stream()
  412. .filter(streamPushItem-> streamPushItem.getGbId() != null)
  413. .collect(Collectors.toList());
  414. // 存储数据到gb_stream表, id会返回到streamPushItemForGbStream里
  415. if (streamPushItemForGbStream.size() > 0) {
  416. gbStreamMapper.batchAdd(streamPushItemForGbStream);
  417. }
  418. // 去除没有ID也就是没有存储到数据库的数据
  419. List<StreamPushItem> streamPushItemsForPlatform = streamPushItemForGbStream.stream()
  420. .filter(streamPushItem-> streamPushItem.getGbStreamId() != null)
  421. .collect(Collectors.toList());
  422. if (streamPushItemsForPlatform.size() > 0) {
  423. // 获取所有平台,平台和目录信息一般不会特别大量。
  424. List<ParentPlatform> parentPlatformList = parentPlatformMapper.getParentPlatformList();
  425. Map<String, Map<String, PlatformCatalog>> platformInfoMap = new HashMap<>();
  426. if (parentPlatformList.size() == 0) {
  427. return;
  428. }
  429. for (ParentPlatform platform : parentPlatformList) {
  430. Map<String, PlatformCatalog> catalogMap = new HashMap<>();
  431. // 创建根节点
  432. PlatformCatalog platformCatalog = new PlatformCatalog();
  433. platformCatalog.setId(platform.getServerGBId());
  434. catalogMap.put(platform.getServerGBId(), platformCatalog);
  435. // 查询所有节点信息
  436. List<PlatformCatalog> platformCatalogs = platformCatalogMapper.selectByPlatForm(platform.getServerGBId());
  437. if (platformCatalogs.size() > 0) {
  438. for (PlatformCatalog catalog : platformCatalogs) {
  439. catalogMap.put(catalog.getId(), catalog);
  440. }
  441. }
  442. platformInfoMap.put(platform.getServerGBId(), catalogMap);
  443. }
  444. List<StreamPushItem> streamPushItemListFroPlatform = new ArrayList<>();
  445. Map<String, List<GbStream>> platformForEvent = new HashMap<>();
  446. // 遍历存储结果,查找app+Stream->platformId+catalogId的对应关系,然后执行批量写入
  447. for (StreamPushItem streamPushItem : streamPushItemsForPlatform) {
  448. List<String[]> platFormInfoList = streamPushItemsForAll.get(streamPushItem.getApp() + streamPushItem.getStream());
  449. if (platFormInfoList != null && platFormInfoList.size() > 0) {
  450. for (String[] platFormInfoArray : platFormInfoList) {
  451. StreamPushItem streamPushItemForPlatform = new StreamPushItem();
  452. streamPushItemForPlatform.setGbStreamId(streamPushItem.getGbStreamId());
  453. if (platFormInfoArray.length > 0) {
  454. // 数组 platFormInfoArray 0 为平台ID。 1为目录ID
  455. // 不存在这个平台,则忽略导入此关联关系
  456. if (platformInfoMap.get(platFormInfoArray[0]) == null
  457. || platformInfoMap.get(platFormInfoArray[0]).get(platFormInfoArray[1]) == null) {
  458. logger.info("导入数据时不存在平台或目录{}/{},已导入未分配", platFormInfoArray[0], platFormInfoArray[1] );
  459. continue;
  460. }
  461. streamPushItemForPlatform.setPlatformId(platFormInfoArray[0]);
  462. List<GbStream> gbStreamList = platformForEvent.get(platFormInfoArray[0]);
  463. if (gbStreamList == null) {
  464. gbStreamList = new ArrayList<>();
  465. platformForEvent.put(platFormInfoArray[0], gbStreamList);
  466. }
  467. // 为发送通知整理数据
  468. streamPushItemForPlatform.setName(streamPushItem.getName());
  469. streamPushItemForPlatform.setApp(streamPushItem.getApp());
  470. streamPushItemForPlatform.setStream(streamPushItem.getStream());
  471. streamPushItemForPlatform.setGbId(streamPushItem.getGbId());
  472. gbStreamList.add(streamPushItemForPlatform);
  473. }
  474. if (platFormInfoArray.length > 1) {
  475. streamPushItemForPlatform.setCatalogId(platFormInfoArray[1]);
  476. }
  477. streamPushItemListFroPlatform.add(streamPushItemForPlatform);
  478. }
  479. }
  480. }
  481. if (!streamPushItemListFroPlatform.isEmpty()) {
  482. platformGbStreamMapper.batchAdd(streamPushItemListFroPlatform);
  483. // 发送通知
  484. for (String platformId : platformForEvent.keySet()) {
  485. eventPublisher.catalogEventPublishForStream(
  486. platformId, platformForEvent.get(platformId), CatalogEvent.ADD);
  487. }
  488. }
  489. }
  490. }
  491. @Override
  492. public boolean batchStop(List<GbStream> gbStreams) {
  493. if (gbStreams == null || gbStreams.size() == 0) {
  494. return false;
  495. }
  496. gbStreamService.sendCatalogMsgs(gbStreams, CatalogEvent.DEL);
  497. platformGbStreamMapper.delByGbStreams(gbStreams);
  498. gbStreamMapper.batchDelForGbStream(gbStreams);
  499. int delStream = streamPushMapper.delAllForGbStream(gbStreams);
  500. if (delStream > 0) {
  501. for (GbStream gbStream : gbStreams) {
  502. MediaServer mediaServerItem = mediaServerService.getOne(gbStream.getMediaServerId());
  503. mediaServerService.closeStreams(mediaServerItem, gbStream.getApp(), gbStream.getStream());
  504. }
  505. }
  506. return true;
  507. }
  508. @Override
  509. public void allStreamOffline() {
  510. List<GbStream> onlinePushers = streamPushMapper.getOnlinePusherForGb();
  511. if (onlinePushers.size() == 0) {
  512. return;
  513. }
  514. streamPushMapper.setAllStreamOffline();
  515. // 发送通知
  516. eventPublisher.catalogEventPublishForStream(null, onlinePushers, CatalogEvent.OFF);
  517. }
  518. @Override
  519. public void offline(List<StreamPushItemFromRedis> offlineStreams) {
  520. // 更新部分设备离线
  521. List<GbStream> onlinePushers = streamPushMapper.getOnlinePusherForGbInList(offlineStreams);
  522. streamPushMapper.offline(offlineStreams);
  523. // 发送通知
  524. eventPublisher.catalogEventPublishForStream(null, onlinePushers, CatalogEvent.OFF);
  525. }
  526. @Override
  527. public void online(List<StreamPushItemFromRedis> onlineStreams) {
  528. // 更新部分设备上线streamPushService
  529. List<GbStream> onlinePushers = streamPushMapper.getOfflinePusherForGbInList(onlineStreams);
  530. streamPushMapper.online(onlineStreams);
  531. // 发送通知
  532. eventPublisher.catalogEventPublishForStream(null, onlinePushers, CatalogEvent.ON);
  533. }
  534. @Override
  535. public boolean add(StreamPushItem stream) {
  536. stream.setUpdateTime(DateUtil.getNow());
  537. stream.setCreateTime(DateUtil.getNow());
  538. stream.setServerId(userSetting.getServerId());
  539. stream.setMediaServerId(mediaConfig.getId());
  540. stream.setSelf(true);
  541. stream.setPushIng(true);
  542. // 放在事务内执行
  543. boolean result = false;
  544. TransactionStatus transactionStatus = dataSourceTransactionManager.getTransaction(transactionDefinition);
  545. try {
  546. int addStreamResult = streamPushMapper.add(stream);
  547. if (!ObjectUtils.isEmpty(stream.getGbId())) {
  548. stream.setStreamType("push");
  549. gbStreamMapper.add(stream);
  550. }
  551. dataSourceTransactionManager.commit(transactionStatus);
  552. result = true;
  553. }catch (Exception e) {
  554. logger.error("批量移除流与平台的关系时错误", e);
  555. dataSourceTransactionManager.rollback(transactionStatus);
  556. }
  557. return result;
  558. }
  559. @Override
  560. public List<String> getAllAppAndStream() {
  561. return streamPushMapper.getAllAppAndStream();
  562. }
  563. @Override
  564. public ResourceBaseInfo getOverview() {
  565. int total = streamPushMapper.getAllCount();
  566. int online = streamPushMapper.getAllOnline(userSetting.isUsePushingAsStatus());
  567. return new ResourceBaseInfo(total, online);
  568. }
  569. @Override
  570. public Map<String, StreamPushItem> getAllAppAndStreamMap() {
  571. return streamPushMapper.getAllAppAndStreamMap();
  572. }
  573. }