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 dataStreamAdapterCache = new ConcurrentHashMap<>(2048); /** * 数据有变化适配器队列 */ private LinkedBlockingQueue readQueue = new LinkedBlockingQueue<>(); /** * 模板处理器链 */ private Map 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(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