| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299 |
- package org.ts.ddcs.dataStream;
- import lombok.extern.slf4j.Slf4j;
- import org.glassfish.jersey.internal.guava.ThreadFactoryBuilder;
- import org.rocksdb.RocksDB;
- import org.ts.ddcs.common.enums.AgreementEnum;
- import org.ts.ddcs.constant.BusinessConstant;
- import org.ts.ddcs.datatemplate.DataTemplateManageHolder;
- import org.ts.ddcs.datatemplate.Template;
- import org.ts.ddcs.itemdatatemplate.ItemTemplate;
- import java.util.List;
- import java.util.Map;
- import java.util.Set;
- import java.util.concurrent.*;
- /**
- * 数据流主模块
- *
- * @author dylan
- */
- @Slf4j
- public class DataStreamManager {
- private RocksDB rocksDB;
- /**
- * 数据流工作线程数量
- */
- private static final int STREAM_THREAD_NUM = 4;
- private int taskIdelCount = 0;
- /**
- * 连接适配器
- */
- private Map<String, DataStreamAdapter> dataStreamAdapterCache = new ConcurrentHashMap<>(2048);
- /**
- * 数据有变化适配器队列
- */
- private LinkedBlockingQueue<String> readQueue = new LinkedBlockingQueue<>();
- /**
- * 模板处理器链
- */
- private Map<Integer, DataStreamLineManager> streamLineCache = new ConcurrentHashMap<>(DataStreamManager.STREAM_THREAD_NUM);
- /**
- * 线程池
- **/
- private static ThreadFactory batchThreadFactory = new ThreadFactoryBuilder().setNameFormat("batch-thread-pool-%d").build();
- private static ExecutorService batchThreadPool = new ThreadPoolExecutor(4, 4,
- 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(1024), batchThreadFactory, new ThreadPoolExecutor.AbortPolicy());
- /**
- * 任务定时器
- */
- private ScheduledExecutorService scheduleExecutor = Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "batch-stream-task-" + r.hashCode()));
- /**
- * 创建数据流处理链
- *
- * @throws IllegalAccessException
- * @throws InstantiationException
- * @throws ClassNotFoundException
- */
- private void createDataStreamLine(String templates) throws IllegalAccessException, InstantiationException, ClassNotFoundException {
- for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
- DataStreamLineManager lineManager = new DataStreamLineManager();
- String[] list = templates.split(",");
- if (list.length > 0) {
- for (String templateName : list) {
- if (AgreementEnum.AGREEMENT_SW.getName().equals(templateName)) {
- List<Template> templateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneTemplateLine(AgreementEnum.AGREEMENT_SW.getName());
- if (null != templateList && templateList.size() > 0) {
- for (Template template : templateList) {
- if (BusinessConstant.TEMPLATE_TYPE_PROCESSOR.equals(template.getType())) {
- lineManager.addProcessor(AgreementEnum.AGREEMENT_SW.getName(), template);
- } else if (BusinessConstant.TEMPLATE_TYPE_SINK.equals(template.getType())) {
- lineManager.addSink(AgreementEnum.AGREEMENT_SW.getName(), template);
- }
- }
- }
- List<ItemTemplate> itemTemplateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneItemTemplateLine(AgreementEnum.AGREEMENT_SW.getName());
- if (null != itemTemplateList && itemTemplateList.size() > 0) {
- for (ItemTemplate itemTemplate : itemTemplateList) {
- lineManager.addItemTemplates(AgreementEnum.AGREEMENT_SW.getName(), itemTemplate);
- }
- }
- } else if (AgreementEnum.AGREEMENT_SZY.getName().equals(templateName)) {
- List<Template> templateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneTemplateLine(AgreementEnum.AGREEMENT_SZY.getName());
- if (null != templateList && templateList.size() > 0) {
- for (Template template : templateList) {
- if (BusinessConstant.TEMPLATE_TYPE_PROCESSOR.equals(template.getType())) {
- lineManager.addProcessor(AgreementEnum.AGREEMENT_SZY.getName(), template);
- } else if (BusinessConstant.TEMPLATE_TYPE_SINK.equals(template.getType())) {
- lineManager.addSink(AgreementEnum.AGREEMENT_SZY.getName(), template);
- }
- }
- }
- List<ItemTemplate> itemTemplateList = DataTemplateManageHolder.newInstance().dataTemplateManager.cloneItemTemplateLine(AgreementEnum.AGREEMENT_SZY.getName());
- if (null != itemTemplateList && itemTemplateList.size() > 0) {
- for (ItemTemplate itemTemplate : itemTemplateList) {
- lineManager.addItemTemplates(AgreementEnum.AGREEMENT_SZY.getName(), itemTemplate);
- }
- }
- }
- }
- }
- streamLineCache.put(i, lineManager);
- }
- }
- /**
- * 数据流管理器初始化
- */
- public void init(String templates,RocksDB rocksDB) throws IllegalAccessException, ClassNotFoundException, InstantiationException {
- try {
- this.rocksDB = rocksDB;
- //初始化数据流处理器
- this.createDataStreamLine(templates);
- //任务定时器
- this.scheduleExecutor.scheduleWithFixedDelay(() -> {
- try {
- boolean idel = false;
- boolean hasTask = false;
- for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
- DataStreamLineManager lineManager = streamLineCache.get(i);
- if (!lineManager.lock()) {
- //有空闲处理器
- String key = readQueue.poll();
- if (null != key) {
- DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
- if (null != dataStreamAdapter) {
- if (!dataStreamAdapter.isLock()) {
- //数据流适配器空闲
- if (dataStreamAdapter.remaining()) {
- lineManager.lock(true);
- dataStreamAdapter.setLock(true);
- //创建数据处理任务
- long tm = System.currentTimeMillis();
- BatchDataStreamTask task = new BatchDataStreamTask(tm, lineManager, dataStreamAdapter);
- FutureTask<Integer> futureTask = new FutureTask<>(task);
- batchThreadPool.execute(futureTask);
- hasTask = true;
- }
- }
- } else {
- log.info("找不到适配器 {}", key);
- }
- }
- idel = true;
- } else {
- log.info("数据流处理器已被锁住 {}", i);
- }
- }
- if (!idel) {
- log.info("任务处理器性能故障");
- }
- if (!hasTask) {
- taskIdelCount += 1;
- } else {
- taskIdelCount = 0;
- }
- //log.info("task idel count {}",taskIdelCount);
- if (taskIdelCount >= 10) {
- //没有任务
- Set<String> keySet = dataStreamAdapterCache.keySet();
- for (String key : keySet) {
- DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
- if (!dataStreamAdapter.isLock()) {
- //数据流适配器空闲
- if (dataStreamAdapter.remaining()) {
- idel = false;
- //有待处理数据
- for (int i = 1; i <= STREAM_THREAD_NUM; i++) {
- DataStreamLineManager lineManager = streamLineCache.get(i);
- if (!lineManager.lock()) {
- //有空闲处理器
- dataStreamAdapter.setLock(true);
- lineManager.lock(true);
- //创建数据处理任务
- long tm = System.currentTimeMillis();
- BatchDataStreamTask task = new BatchDataStreamTask(tm, lineManager, dataStreamAdapter);
- FutureTask<Integer> futureTask = new FutureTask<>(task);
- batchThreadPool.execute(futureTask);
- idel = true;
- break;
- } else {
- log.info("数据流处理器已被锁住 {}", i);
- }
- }
- if (!idel) {
- log.info("任务处理器性能故障");
- break;
- }
- }
- }
- }
- taskIdelCount = 0;
- }
- } catch (Exception e) {
- log.error(this.getClass().getName(), e);
- }
- }, 0, 10, TimeUnit.MILLISECONDS);
- } catch (Exception e) {
- log.error(this.getClass().getName(), e);
- throw e;
- }
- }
- private class BatchDataStreamTask implements Callable<Integer> {
- private long batchTime;
- private DataStreamAdapter dataStreamAdapter;
- private DataStreamLineManager dataStreamLineManager;
- BatchDataStreamTask(long batchTime, DataStreamLineManager dataStreamLineManager, DataStreamAdapter dataStreamAdapter) {
- this.batchTime = batchTime;
- this.dataStreamLineManager = dataStreamLineManager;
- this.dataStreamAdapter = dataStreamAdapter;
- }
- @Override
- public Integer call() {
- try {
- dataStreamLineManager.stream(dataStreamAdapter);
- dataStreamAdapter.setLock(false);
- dataStreamLineManager.lock(false);
- // long ut = System.currentTimeMillis() - batchTm;
- // LogHelper.info("time " + ut);
- // if (ut > 10) {
- // LogHelper.info("task time too long !!!!!!!!!!!!!!!");
- // }
- } catch (Exception e) {
- log.error(this.getClass().getName(), e);
- }
- return 0;
- }
- }
- /**
- * 注册连接
- *
- * @param adapter
- */
- public void regConnect(DataStreamAdapter adapter) {
- dataStreamAdapterCache.put(adapter.getAdapterKey(), adapter);
- }
- /**
- * 关闭连接
- *
- * @param adapter
- */
- public void closeConnect(DataStreamAdapter adapter) {
- if (dataStreamAdapterCache.containsKey(adapter.getAdapterKey())) {
- if (!adapter.isLock()) {
- byte[] block = adapter.getBlock();
- if (null == block || block.length == 0) {
- dataStreamAdapterCache.remove(adapter.getAdapterKey());
- }
- }
- }
- }
- public void readEvent(DataStreamAdapter adapter) {
- readQueue.offer(adapter.getAdapterKey());
- }
- public void responseEvent(String key, byte[] respBlock) {
- if (this.dataStreamAdapterCache.containsKey(key)) {
- DataStreamAdapter dataStreamAdapter = this.dataStreamAdapterCache.get(key);
- if (!dataStreamAdapter.isClose()) {
- dataStreamAdapter.output(respBlock);
- } else {
- log.info("运行错误,链接已经提前关闭: {}", key);
- }
- }
- }
- public DataStreamAdapter getAdapter(String rtuCode){
- Set<String> keySet = dataStreamAdapterCache.keySet();
- for (String key : keySet) {
- DataStreamAdapter dataStreamAdapter = dataStreamAdapterCache.get(key);
- String stcd =dataStreamAdapter.getRtuCode();
- if (null != stcd && stcd.equals(rtuCode)){
- return dataStreamAdapter;
- }
- }
- return null;
- }
- }
|