DataStreamManager.java 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299
  1. package org.ts.ddcs.dataStream;
  2. import lombok.extern.slf4j.Slf4j;
  3. import org.glassfish.jersey.internal.guava.ThreadFactoryBuilder;
  4. import org.rocksdb.RocksDB;
  5. import org.ts.ddcs.common.enums.AgreementEnum;
  6. import org.ts.ddcs.constant.BusinessConstant;
  7. import org.ts.ddcs.datatemplate.DataTemplateManageHolder;
  8. import org.ts.ddcs.datatemplate.Template;
  9. import org.ts.ddcs.itemdatatemplate.ItemTemplate;
  10. import java.util.List;
  11. import java.util.Map;
  12. import java.util.Set;
  13. import java.util.concurrent.*;
  14. /**
  15. * 数据流主模块
  16. *
  17. * @author dylan
  18. */
  19. @Slf4j
  20. public class DataStreamManager {
  21. private RocksDB rocksDB;
  22. /**
  23. * 数据流工作线程数量
  24. */
  25. private static final int STREAM_THREAD_NUM = 4;
  26. private int taskIdelCount = 0;
  27. /**
  28. * 连接适配器
  29. */
  30. private Map<String, DataStreamAdapter> dataStreamAdapterCache = new ConcurrentHashMap<>(2048);
  31. /**
  32. * 数据有变化适配器队列
  33. */
  34. private LinkedBlockingQueue<String> readQueue = new LinkedBlockingQueue<>();
  35. /**
  36. * 模板处理器链
  37. */
  38. private Map<Integer, DataStreamLineManager> streamLineCache = new ConcurrentHashMap<>(DataStreamManager.STREAM_THREAD_NUM);
  39. /**
  40. * 线程池
  41. **/
  42. private static ThreadFactory batchThreadFactory = new ThreadFactoryBuilder().setNameFormat("batch-thread-pool-%d").build();
  43. private static ExecutorService batchThreadPool = new ThreadPoolExecutor(4, 4,
  44. 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(1024), batchThreadFactory, new ThreadPoolExecutor.AbortPolicy());
  45. /**
  46. * 任务定时器
  47. */
  48. private ScheduledExecutorService scheduleExecutor = Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "batch-stream-task-" + r.hashCode()));
  49. /**
  50. * 创建数据流处理链
  51. *
  52. * @throws IllegalAccessException
  53. * @throws InstantiationException
  54. * @throws ClassNotFoundException
  55. */
  56. private void createDataStreamLine(String templates) throws IllegalAccessException, InstantiationException, ClassNotFoundException {
  57. for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
  58. DataStreamLineManager lineManager = new DataStreamLineManager();
  59. String[] list = templates.split(",");
  60. if (list.length > 0) {
  61. for (String templateName : list) {
  62. if (AgreementEnum.AGREEMENT_SW.getName().equals(templateName)) {
  63. List<Template> templateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneTemplateLine(AgreementEnum.AGREEMENT_SW.getName());
  64. if (null != templateList && templateList.size() > 0) {
  65. for (Template template : templateList) {
  66. if (BusinessConstant.TEMPLATE_TYPE_PROCESSOR.equals(template.getType())) {
  67. lineManager.addProcessor(AgreementEnum.AGREEMENT_SW.getName(), template);
  68. } else if (BusinessConstant.TEMPLATE_TYPE_SINK.equals(template.getType())) {
  69. lineManager.addSink(AgreementEnum.AGREEMENT_SW.getName(), template);
  70. }
  71. }
  72. }
  73. List<ItemTemplate> itemTemplateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneItemTemplateLine(AgreementEnum.AGREEMENT_SW.getName());
  74. if (null != itemTemplateList && itemTemplateList.size() > 0) {
  75. for (ItemTemplate itemTemplate : itemTemplateList) {
  76. lineManager.addItemTemplates(AgreementEnum.AGREEMENT_SW.getName(), itemTemplate);
  77. }
  78. }
  79. } else if (AgreementEnum.AGREEMENT_SZY.getName().equals(templateName)) {
  80. List<Template> templateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneTemplateLine(AgreementEnum.AGREEMENT_SZY.getName());
  81. if (null != templateList && templateList.size() > 0) {
  82. for (Template template : templateList) {
  83. if (BusinessConstant.TEMPLATE_TYPE_PROCESSOR.equals(template.getType())) {
  84. lineManager.addProcessor(AgreementEnum.AGREEMENT_SZY.getName(), template);
  85. } else if (BusinessConstant.TEMPLATE_TYPE_SINK.equals(template.getType())) {
  86. lineManager.addSink(AgreementEnum.AGREEMENT_SZY.getName(), template);
  87. }
  88. }
  89. }
  90. List<ItemTemplate> itemTemplateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneItemTemplateLine(AgreementEnum.AGREEMENT_SZY.getName());
  91. if (null != itemTemplateList && itemTemplateList.size() > 0) {
  92. for (ItemTemplate itemTemplate : itemTemplateList) {
  93. lineManager.addItemTemplates(AgreementEnum.AGREEMENT_SZY.getName(), itemTemplate);
  94. }
  95. }
  96. }
  97. }
  98. }
  99. streamLineCache.put(i, lineManager);
  100. }
  101. }
  102. /**
  103. * 数据流管理器初始化
  104. */
  105. public void init(String templates,RocksDB rocksDB) throws IllegalAccessException, ClassNotFoundException, InstantiationException {
  106. try {
  107. this.rocksDB = rocksDB;
  108. //初始化数据流处理器
  109. this.createDataStreamLine(templates);
  110. //任务定时器
  111. this.scheduleExecutor.scheduleWithFixedDelay(() -> {
  112. try {
  113. boolean idel = false;
  114. boolean hasTask = false;
  115. for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
  116. DataStreamLineManager lineManager = streamLineCache.get(i);
  117. if (!lineManager.lock()) {
  118. //有空闲处理器
  119. String key = readQueue.poll();
  120. if (null != key) {
  121. DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
  122. if (null != dataStreamAdapter) {
  123. if (!dataStreamAdapter.isLock()) {
  124. //数据流适配器空闲
  125. if (dataStreamAdapter.remaining()) {
  126. lineManager.lock(true);
  127. dataStreamAdapter.setLock(true);
  128. //创建数据处理任务
  129. long tm = System.currentTimeMillis();
  130. BatchDataStreamTask task = new BatchDataStreamTask(tm, lineManager, dataStreamAdapter);
  131. FutureTask<Integer> futureTask = new FutureTask<>(task);
  132. batchThreadPool.execute(futureTask);
  133. hasTask = true;
  134. }
  135. }
  136. } else {
  137. log.info("找不到适配器 {}", key);
  138. }
  139. }
  140. idel = true;
  141. } else {
  142. log.info("数据流处理器已被锁住 {}", i);
  143. }
  144. }
  145. if (!idel) {
  146. log.info("任务处理器性能故障");
  147. }
  148. if (!hasTask) {
  149. taskIdelCount += 1;
  150. } else {
  151. taskIdelCount = 0;
  152. }
  153. //log.info("task idel count {}",taskIdelCount);
  154. if (taskIdelCount >= 10) {
  155. //没有任务
  156. Set<String> keySet = dataStreamAdapterCache.keySet();
  157. for (String key : keySet) {
  158. DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
  159. if (!dataStreamAdapter.isLock()) {
  160. //数据流适配器空闲
  161. if (dataStreamAdapter.remaining()) {
  162. idel = false;
  163. //有待处理数据
  164. for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
  165. DataStreamLineManager lineManager = streamLineCache.get(i);
  166. if (!lineManager.lock()) {
  167. //有空闲处理器
  168. dataStreamAdapter.setLock(true);
  169. lineManager.lock(true);
  170. //创建数据处理任务
  171. long tm = System.currentTimeMillis();
  172. BatchDataStreamTask task = new BatchDataStreamTask(tm, lineManager, dataStreamAdapter);
  173. FutureTask<Integer> futureTask = new FutureTask<>(task);
  174. batchThreadPool.execute(futureTask);
  175. idel = true;
  176. break;
  177. } else {
  178. log.info("数据流处理器已被锁住 {}", i);
  179. }
  180. }
  181. if (!idel) {
  182. log.info("任务处理器性能故障");
  183. break;
  184. }
  185. }
  186. }
  187. }
  188. taskIdelCount = 0;
  189. }
  190. } catch (Exception e) {
  191. log.error(this.getClass().getName(), e);
  192. }
  193. }, 0, 10, TimeUnit.MILLISECONDS);
  194. } catch (Exception e) {
  195. log.error(this.getClass().getName(), e);
  196. throw e;
  197. }
  198. }
  199. private class BatchDataStreamTask implements Callable<Integer> {
  200. private long batchTime;
  201. private DataStreamAdapter dataStreamAdapter;
  202. private DataStreamLineManager dataStreamLineManager;
  203. BatchDataStreamTask(long batchTime, DataStreamLineManager dataStreamLineManager, DataStreamAdapter dataStreamAdapter) {
  204. this.batchTime = batchTime;
  205. this.dataStreamLineManager = dataStreamLineManager;
  206. this.dataStreamAdapter = dataStreamAdapter;
  207. }
  208. @Override
  209. public Integer call() {
  210. try {
  211. dataStreamLineManager.stream(dataStreamAdapter);
  212. dataStreamAdapter.setLock(false);
  213. dataStreamLineManager.lock(false);
  214. // long ut = System.currentTimeMillis() - batchTm;
  215. // LogHelper.info("time " + ut);
  216. // if (ut > 10) {
  217. // LogHelper.info("task time too long !!!!!!!!!!!!!!!");
  218. // }
  219. } catch (Exception e) {
  220. log.error(this.getClass().getName(), e);
  221. }
  222. return 0;
  223. }
  224. }
  225. /**
  226. * 注册连接
  227. *
  228. * @param adapter
  229. */
  230. public void regConnect(DataStreamAdapter adapter) {
  231. dataStreamAdapterCache.put(adapter.getAdapterKey(), adapter);
  232. }
  233. /**
  234. * 关闭连接
  235. *
  236. * @param adapter
  237. */
  238. public void closeConnect(DataStreamAdapter adapter) {
  239. if (dataStreamAdapterCache.containsKey(adapter.getAdapterKey())) {
  240. if (!adapter.isLock()) {
  241. byte[] block = adapter.getBlock();
  242. if (null == block || block.length == 0) {
  243. dataStreamAdapterCache.remove(adapter.getAdapterKey());
  244. }
  245. }
  246. }
  247. }
  248. public void readEvent(DataStreamAdapter adapter) {
  249. readQueue.offer(adapter.getAdapterKey());
  250. }
  251. public void responseEvent(String key, byte[] respBlock) {
  252. if (this.dataStreamAdapterCache.containsKey(key)) {
  253. DataStreamAdapter dataStreamAdapter = this.dataStreamAdapterCache.get(key);
  254. if (!dataStreamAdapter.isClose()) {
  255. dataStreamAdapter.output(respBlock);
  256. } else {
  257. log.info("运行错误,链接已经提前关闭: {}", key);
  258. }
  259. }
  260. }
  261. public DataStreamAdapter getAdapter(String rtuCode){
  262. Set<String> keySet = dataStreamAdapterCache.keySet();
  263. for (String key : keySet) {
  264. DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
  265. String stcd =dataStreamAdapter.getRtuCode();
  266. if (null != stcd && stcd.equals(rtuCode)){
  267. return dataStreamAdapter;
  268. }
  269. }
  270. return null;
  271. }
  272. }