MediaService.java 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170
  1. package com.zj.service;
  2. import java.io.ByteArrayOutputStream;
  3. import java.io.IOException;
  4. import java.util.HashMap;
  5. import java.util.Map;
  6. import javax.servlet.http.HttpServletResponse;
  7. import org.springframework.stereotype.Service;
  8. import org.springframework.web.socket.BinaryMessage;
  9. import org.springframework.web.socket.WebSocketSession;
  10. import com.zj.thread.ProcessThread;
  11. import cn.hutool.core.thread.ThreadUtil;
  12. /**
  13. * 媒体服务
  14. * @author ZJ
  15. *
  16. */
  17. @Service
  18. public class MediaService {
  19. // 缓存线程
  20. public static Map<String, ProcessThread> map = new HashMap<String, ProcessThread>();
  21. /**
  22. *
  23. * @param url 源地址
  24. * @param id 源地址唯一标识(表示同一个媒体)
  25. */
  26. public void playForHttp(String input, String id, HttpServletResponse response) {
  27. ProcessThread processThread = map.get(id);
  28. //新增的媒体需要进行推流初始化
  29. if (processThread == null) {
  30. processThread = new ProcessThread(input);
  31. //初始化推拉流
  32. map.put(id, processThread);
  33. ThreadUtil.execute(processThread);
  34. }
  35. //创建客户端的输出流
  36. ByteArrayOutputStream byteArrayOutputStream = processThread.addClient();
  37. //发送头部
  38. sendHeaderForHttp(response, processThread);
  39. //发送数据
  40. sendAVDataForHttp(response, byteArrayOutputStream);
  41. }
  42. public void playForWs(String input, String id, WebSocketSession session) {
  43. ProcessThread processThread = map.get(id);
  44. //新增的媒体需要进行推流初始化
  45. if (processThread == null) {
  46. processThread = new ProcessThread(input);
  47. //初始化推拉流
  48. map.put(id, processThread);
  49. ThreadUtil.execute(processThread);
  50. }
  51. //创建客户端的输出流
  52. ByteArrayOutputStream byteArrayOutputStream = processThread.addClient();
  53. //发送头部
  54. sendHeaderForWs(session, processThread);
  55. //发送数据
  56. sendAVDataForWs(session, byteArrayOutputStream);
  57. }
  58. /**
  59. * 发送FLV header
  60. * @param response
  61. * @param stream
  62. */
  63. private void sendHeaderForHttp(HttpServletResponse response, ProcessThread processThread) {
  64. try {
  65. //最多等三分钟,如果没有header则认为没取到视频,发送header后续要优化
  66. for (int i = 0; i < 1200; i++) {
  67. if (processThread.getHeader() != null) {
  68. response.getOutputStream().write(processThread.getHeader());
  69. break;
  70. }
  71. Thread.sleep(100);
  72. }
  73. /*
  74. * 这里后续还要修改,如果没获取到header怎么和前段交互,后端线程怎么处理
  75. */
  76. } catch (IOException e) {
  77. e.printStackTrace();
  78. } catch (InterruptedException e) {
  79. }
  80. }
  81. /**
  82. * 发送视频数据
  83. * @param response
  84. * @param outData
  85. */
  86. private void sendAVDataForHttp(HttpServletResponse response, ByteArrayOutputStream outData) {
  87. try {
  88. while (true) {
  89. if (outData.size() > 0) {
  90. response.getOutputStream().write(outData.toByteArray());
  91. outData.reset();
  92. } else {
  93. Thread.sleep(100);
  94. }
  95. }
  96. } catch (IOException e) {
  97. e.printStackTrace();
  98. } catch (InterruptedException e) {
  99. }
  100. }
  101. /**
  102. * 发送FLV header
  103. * @param response
  104. * @param stream
  105. */
  106. private void sendHeaderForWs(WebSocketSession session, ProcessThread processThread) {
  107. try {
  108. //最多等三分钟,如果没有header则认为没取到视频,发送header后续要优化
  109. for (int i = 0; i < 1200; i++) {
  110. if (processThread.getHeader() != null) {
  111. session.sendMessage(new BinaryMessage(processThread.getHeader()));
  112. break;
  113. }
  114. Thread.sleep(100);
  115. }
  116. /*
  117. * 这里后续还要修改,如果没获取到header怎么和前段交互,后端线程怎么处理
  118. */
  119. } catch (IOException e) {
  120. e.printStackTrace();
  121. } catch (InterruptedException e) {
  122. }
  123. }
  124. /**
  125. * 发送视频数据
  126. * @param response
  127. * @param outData
  128. */
  129. private void sendAVDataForWs(WebSocketSession session, ByteArrayOutputStream outData) {
  130. try {
  131. while (true) {
  132. if (outData.size() > 0) {
  133. session.sendMessage(new BinaryMessage(outData.toByteArray()));
  134. outData.reset();
  135. } else {
  136. Thread.sleep(100);
  137. }
  138. }
  139. } catch (IOException e) {
  140. e.printStackTrace();
  141. } catch (InterruptedException e) {
  142. }
  143. }
  144. }