package com.robustel.rlink.device.service.impl;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Date;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import javax.annotation.PostConstruct;
import javax.jms.Destination;
import org.apache.activemq.command.ActiveMQTopic;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.Sort;
import org.springframework.data.domain.Sort.Direction;
import org.springframework.data.domain.Sort.Order;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.data.mongodb.core.query.Update;
import com.robustel.iot.data.share.entity.DeviceModuleData;
import com.robustel.pl.util.constant.DeviceConstantTemplate;
import com.robustel.pl.util.thread.QunarThreadPoolExecutor;
import com.robustel.pl.util.utils.AppUtils;
import com.robustel.pl.util.utils.JsonUtil;
import com.robustel.rlink.control.entity.DeviceCommandRequest;
import com.robustel.rlink.control.enums.CommandType;
import com.robustel.rlink.device.entity.Device;
import com.robustel.rlink.device.enums.DeviceStatusEnum;
import com.robustel.rlink.device.service.RegisterResponseSenderService;
import com.robustel.rlink.device.vo.ModelData;
/**
* 定时处理任务--
* @author jfn
*
*/
public class DeviceOverTimeHandler {
public static DeviceOverTimeHandler handler;
static Logger logger = LoggerFactory.getLogger(DeviceOverTimeHandler.class);
@Autowired
private MongoTemplate mongoTemplate;
@Autowired
private RegisterResponseSenderService registerResponseSenderService;
@PostConstruct
public void init() {
handler = this;
handler.mongoTemplate = this.mongoTemplate;
handler.registerResponseSenderService = this.registerResponseSenderService;
}
static String mdRealTime = DeviceConstantTemplate.module_data_real_time_set;
static Integer betch = 50;
static Integer threadNum = 30;
public static final String CliTopicPerfix="sys_cli.";
public static final String CmdTopicPerfix="sys_ctrl.";
private static List<ModelData> operaSet = Collections.synchronizedList(new ArrayList<ModelData>(betch));
static CountDownLatch startSignal = new CountDownLatch(1);
static CountDownLatch doneSignal = new CountDownLatch(5);
private static ExecutorService pool = Executors.newFixedThreadPool(1024);
private LinkedBlockingQueue<Runnable> queue = new LinkedBlockingQueue<Runnable>();
private QunarThreadPoolExecutor qunarThreadPoolExecutor = new QunarThreadPoolExecutor(50, 200, 5, TimeUnit.MINUTES, queue);
private static Long deviceCount = 0L;
private static Query query = null;
private static volatile Integer num = 0;
public void overTimeProcess(){
logger.info("模块数据下发定时任务启动..."+handler.mongoTemplate);
makeThreadPools();
}
public void makeThreadPools(){
logger.info("线程池初始化...");
query = new Query();
query.addCriteria(Criteria.where("flag").is(-1));
deviceCount = handler.mongoTemplate.count(query,ModelData.class, mdRealTime);
while(deviceCount -betch*num > 0){
query = new Query();
query = query.skip(betch*num).limit(betch);
List<ModelData> data = handler.mongoTemplate.find(query,ModelData.class, mdRealTime);
logger.info("查询数据 ...{}",data);
if(AppUtils.isNotBlank(data)){
Worker worker = new Worker(data);
qunarThreadPoolExecutor.execute(worker);
}
num++;
}
}
public static void addTasks(){
while(operaSet.size()<=0 && deviceCount -betch*num > 0){
query = query.skip(betch*num).limit(betch);
List<ModelData> data = handler.mongoTemplate.find(query,ModelData.class, mdRealTime);
if(AppUtils.isBlank(data)){
/**for(Thread t : pools){
t.stop();
}
pools.clear();**/
pool.shutdown();
logger.info("线程池退出....");
}
operaSet.addAll(data);
num++;
}
}
public static void doWork(List<ModelData> operaSet){
logger.info("working......");
if(AppUtils.isBlank(operaSet)||operaSet.size()<=0){
logger.info("待处理数据为空,线程退出...");
return;
}
logger.info("待处理模块数据 {}",operaSet);
for(ModelData md : operaSet){
logger.info("dowork object {}" ,md);
sendModuleData(md);
}
}
private static void sendModuleData(ModelData md){
logger.info("modelData is {}",md);
Query query = new Query();
query.addCriteria(Criteria.where("_id").is(md.getSn()));
Device dev = handler.mongoTemplate.findOne(query,Device.class, DeviceConstantTemplate.device_list_collection_name);
ModelData data = handler.mongoTemplate.findOne(query,ModelData.class, DeviceConstantTemplate.module_data_real_time_set);
if(data==null||data.getFlag().equals(1)||dev==null || dev.getDeviceOnLineStatus().equals(DeviceStatusEnum.OFF_LINE.value())){
return;
}
DeviceCommandRequest cmd = new DeviceCommandRequest();
cmd.setSn(md.getSn());
cmd.setTime(new Date().getTime());
cmd.setId(md.getCommId());
cmd.setCmd(CommandType.ConfigModuleData.getType());
cmd.setData(md.getPropers());
Destination topic = new ActiveMQTopic(CmdTopicPerfix+md.getSn());
handler.registerResponseSenderService.sendMQTT(topic, JsonUtil.javaObjToJson(cmd));
Update update = new Update();
update.set("flag", 1);
update.set("sendTime", new Date().getTime());
handler.mongoTemplate.updateFirst(query, update, ModelData.class,mdRealTime);
query.addCriteria(Criteria.where("commId").is(md.getCommId()));
handler.mongoTemplate.updateFirst(query, update, ModelData.class,DeviceConstantTemplate.module_data_set);
}
static class Worker implements Runnable{
List<ModelData> operaSet;
public Worker(List<ModelData> operaSet){
this.operaSet = operaSet;
}
public void run() {
logger.info("worker run function ...");
doWork(operaSet);
}
}
private void deviceOffline(Device dev) {
if(AppUtils.isBlank(dev)){
return;
}
Query query2 = new Query();
Criteria cnd2 = Criteria.where("sn").is(dev.getSn());
query2.addCriteria(cnd2).with(new Sort(new Order(Direction.DESC,"time")));
DeviceModuleData module = handler.mongoTemplate.findOne(query2, DeviceModuleData.class, DeviceConstantTemplate.module_data_real_time);
if(AppUtils.isBlank(module)){
//超时(10分钟)没有上传数据的设备将会被下线
upateStatus(module,dev,dev.getDeviceUpdateTime());
}else{
upateStatus(module,dev,module.getTime());
}
}
private void upateStatus(DeviceModuleData module,Device dev,Long time) {
if(time==null)
return;
if((new Date().getTime()-time)>10*60*1000 && dev.getDeviceOnLineStatus().equals(DeviceStatusEnum.ON_LINE.value())){
Update update = new Update();
update.set("deviceOnLineStatus", DeviceStatusEnum.OFF_LINE.value());
update.set("deviceLastLoginTime", dev.getOnTime());
update.set("deviceUpdateTime", new Date().getTime());
update.set("offTime", new Date().getTime());
if(AppUtils.isNotBlank(dev.getOffTime())){
update.set("deviceLastOffLineTime", dev.getOffTime());
}
Query qUpte = new Query();
qUpte.addCriteria(Criteria.where("_id").is(dev.getSn()));
handler.mongoTemplate.updateFirst(qUpte, update, Device.class,DeviceConstantTemplate.device_list_collection_name);
}
}
public MongoTemplate getMongoTemplate() {
return mongoTemplate;
}
public void setMongoTemplate(MongoTemplate mongoTemplate) {
this.mongoTemplate = mongoTemplate;
}
public RegisterResponseSenderService getRegisterResponseSenderService() {
return registerResponseSenderService;
}
public void setRegisterResponseSenderService(
RegisterResponseSenderService registerResponseSenderService) {
this.registerResponseSenderService = registerResponseSenderService;
}
}
package com.robustel.pl.util.thread;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.RejectedExecutionHandler;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
/**
*
* @author jingfangnan
*
*/
public class QunarThreadPoolExecutor extends ThreadPoolExecutor {
// 记录每个线程执行任务开始时间
private ThreadLocal<Long> start = new ThreadLocal<Long>();
// 记录所有任务完成使用的时间
private AtomicLong totals = new AtomicLong();
// 记录线程池完成的任务数
private AtomicInteger tasks = new AtomicInteger();
public QunarThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit,
BlockingQueue<Runnable> workQueue, ThreadFactory threadFactory, RejectedExecutionHandler handler) {
super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler);
}
public QunarThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit,
BlockingQueue<Runnable> workQueue, RejectedExecutionHandler handler) {
super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, handler);
}
public QunarThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit,
BlockingQueue<Runnable> workQueue, ThreadFactory threadFactory) {
super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory);
}
public QunarThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit,
BlockingQueue<Runnable> workQueue) {
super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
}
/**
* 每个线程在调用run方法之前调用该方法
* */
protected void beforeExecute(Thread t, Runnable r) {
super.beforeExecute(t, r);
start.set(System.currentTimeMillis());
}
/**
* 每个线程在执行完run方法后调用该方法
* */
protected void afterExecute(Runnable r, Throwable t) {
super.afterExecute(r, t);
tasks.incrementAndGet();
totals.addAndGet(System.currentTimeMillis() - start.get());
}
@Override
protected void terminated() {
super.terminated();
System.out.println("完成"+ tasks.get() +"个任务,平均耗时: [" + totals.get() / tasks.get() + "] ms");
}
}
分享到:
相关推荐
SSH+JBoss+多线程的电子宠物系统是一种基于Java技术栈开发的复杂应用程序,它融合了Spring、Struts和Hibernate三大主流框架的优势,利用JBoss应用服务器提供服务,并通过多线程技术来实现高性能和高并发处理。...
JAVA核心API包含Java标准库中的各种类和接口,如集合框架、多线程、网络编程等。深入理解并熟练运用这些API,能够有效地解决问题,提高编程效率。 综上所述,这些文件内容涵盖了从Java基础到高级应用的广泛知识,...
**Java多线程并发** Java提供了丰富的多线程并发工具,如Thread、Runnable、synchronized关键字、Lock接口、Future、ExecutorService等,支持线程同步、互斥、协作和异步计算。 **Spring框架原理** Spring是一个...
Java的关键特性包括自动内存管理(垃圾回收)、多线程支持和丰富的类库。EJB(Enterprise JavaBeans)是Java EE平台的一部分,用于构建可部署到服务器端的企业级应用,分为会话bean、实体bean和消息驱动bean。Spring...
Java的特性如多线程、异常处理、丰富的类库使得它在处理复杂业务场景时表现得游刃有余。 2. **JSP(JavaServer Pages)**:JSP是Java Web开发中的视图层技术,用于生成动态网页内容。开发者可以在JSP页面上混合HTML...
- **效率**:HashMap是非同步的,因此在多线程环境下可能会出现并发问题,但其执行效率通常高于Hashtable。 - **contains方法**:HashMap移除了Hashtable中的`contains`方法,提供了`containsValue`和`containsKey`...
这个合集可能还包括线程池管理工具,如`ThreadPoolUtil`,用于创建和管理线程池,提高多线程环境下的性能和可维护性。还有可能包含网络通信工具类,如`HttpClientUtil`,用于发起HTTP请求,处理网络数据。 总之,...
- **语法特性**:Java是一种面向对象的编程语言,具有静态类型系统和自动内存管理。它的语法简洁且易于理解,同时支持封装、继承和多态等面向对象概念。 - **JVM(Java虚拟机)**:Java程序通过JVM运行,JVM负责...
首先,Java SE是Java的基础,包括语法、数据类型、控制结构、类与对象、接口、异常处理、多线程、集合框架等。在面试中,面试官可能会问及以下问题: 1. Java内存模型:理解栈、堆、方法区的差异,以及垃圾回收机制...
在Java高级面试中,面试官通常会关注候选人在核心Java、多线程、集合框架、JVM内存管理、设计模式、数据库操作、网络编程、异常处理、IO流、Spring框架及其实现原理等方面的知识掌握程度。以下是根据这些关键点展开...
在Java面试中,面试官通常会考察候选人的基础知识、面向对象编程理解、JVM工作原理、集合框架、多线程、异常处理、设计模式以及企业级应用开发的相关技术,如J2EE和Spring框架。 1. **Java高效的原因**:Java的高效...
"部分面试题答案.doc"则可能包含了一些常见的Java和C++面试问题,比如多线程、内存管理、设计模式等。 综上所述,软件工程师在Java和C++领域的学习和实践中,不仅需要掌握编程语言本身,还需要了解相关的框架和工具...
- JDBC(Java Database Connectivity):Java访问数据库的标准接口,允许程序连接、查询和操作数据库。 - EJB(Enterprise JavaBeans):企业级组件,包括会话bean(处理业务逻辑)、实体bean(持久化数据)和消息...
需要注意的是,静态字段和方法在多线程环境下可能需要额外的同步控制,以避免并发问题。 对于Servlet、Filter和Listener,由于它们通常在Web应用启动时由容器实例化,而非由Spring管理,所以也不能直接使用@...
3. **多线程与并发**: - 线程的创建方式:Thread类和Runnable接口。 - 线程同步:synchronized关键字,volatile变量,Lock接口和Condition。 - 死锁、活锁、饥饿现象的理解及避免策略。 - Executors框架的使用...
四、多线程与并发 1. 线程创建:通过Thread类或实现Runnable接口来创建线程,理解Thread的start()方法与run()方法的区别。 2. 线程同步:熟悉synchronized关键字,wait()、notify()和notifyAll()方法,以及Lock接口...
标题提到的问题——“从bean工厂里单例执行方法效率比new对象执行慢很多”,涉及到Java编程中的两种常见对象管理方式:单例模式和直接实例化。这个现象可能让开发者感到困惑,因为通常认为单例模式在性能上具有优势...
21. **多线程**:涉及线程同步、并发控制和死锁问题。 22. **文件加密技术**:加密算法如AES、RSA用于保护数据安全。 23. **软件开发生命周期**:包括需求分析、设计、编码、测试和维护阶段。 24. **路由协议**:如...
- **线程共享**:堆(存储对象实例),方法区(存储类信息、常量、静态变量)。 - **线程私有**:程序计数器(记录执行的指令地址),JVM栈(存储方法调用信息),本地方法栈(为JNI方法服务)。 这些知识涵盖了...