使用队列接收客户端消息

dev
lisai17@sina.com 2020-12-30 14:37:56 +08:00
parent 3bcd5fa438
commit 87d135d723
7 changed files with 79 additions and 23 deletions

View File

@ -15,7 +15,7 @@ socketserver.port=21002
socketio.port=12002 socketio.port=12002
#当前部署本地程序的砂站id #当前部署本地程序的砂站id
current.supermarket_id=8 current.supermarket_id=1
#落杆后,等待上磅的时间 #落杆后,等待上磅的时间
default_scale_wait_time=8000 default_scale_wait_time=8000

View File

@ -1,12 +1,12 @@
## mysql ## mysql
## GRANT ALL PRIVILEGES ON *.* TO 'root'@'192.168.1.119' IDENTIFIED BY 'Local_1' WITH GRANT OPTION; ## GRANT ALL PRIVILEGES ON *.* TO 'root'@'192.168.1.119' IDENTIFIED BY 'Local_1' WITH GRANT OPTION;
#jdbcUrl=jdbc:mysql://rm-wz9wa070076b2uge2ro.mysql.rds.aliyuncs.com:3306/ssjy_xsx_dev?characterEncoding=utf8&useSSL=false&zeroDateTimeBehavior=convertToNull&useInformationSchema=true&serverTimezone=GMT%2B8&autoReconnect=true jdbcUrl=jdbc:mysql://rm-wz9wa070076b2uge2ro.mysql.rds.aliyuncs.com:3306/ssjy_xsx_dev?characterEncoding=utf8&useSSL=false&zeroDateTimeBehavior=convertToNull&useInformationSchema=true&serverTimezone=GMT%2B8&autoReconnect=true
#user=dev_ssjy_xsx user=dev_ssjy_xsx
#password=Ssjy_xs_890 password=Ssjy_xs_890
jdbcUrl=jdbc:mysql://192.168.20.2:3306/ssjy_xsx_dev?characterEncoding=utf8&useSSL=false&zeroDateTimeBehavior=convertToNull&useInformationSchema=true&serverTimezone=GMT%2B8&autoReconnect=true #jdbcUrl=jdbc:mysql://192.168.20.2:3306/ssjy_xsx_dev?characterEncoding=utf8&useSSL=false&zeroDateTimeBehavior=convertToNull&useInformationSchema=true&serverTimezone=GMT%2B8&autoReconnect=true
user=root #user=root
password=Ssjy_xsx_890 #password=Ssjy_xsx_890
# redis # redis
redis.basekey=ssjcgl_xsx_dev redis.basekey=ssjcgl_xsx_dev

View File

@ -5,15 +5,14 @@ import com.cowr.common.view.Result;
import com.cowr.service.ssjygl.main.AuthInterceptor; import com.cowr.service.ssjygl.main.AuthInterceptor;
import com.cowr.service.ssjygl.main.SvrCacheData; import com.cowr.service.ssjygl.main.SvrCacheData;
import com.cowr.service.ssjygl.supermarket.SupermarketSyncService; import com.cowr.service.ssjygl.supermarket.SupermarketSyncService;
import com.cowr.service.ssjygl.synctask.SyncTaskService;
import com.cowr.ssjygl.transportcompany.TransportCompanyService; import com.cowr.ssjygl.transportcompany.TransportCompanyService;
import com.cowr.ssjygl.transprice.TransPriceService; import com.cowr.ssjygl.transprice.TransPriceService;
import com.jfinal.aop.Clear; import com.jfinal.aop.Clear;
import com.jfinal.core.Controller; import com.jfinal.core.Controller;
import com.jfinal.plugin.activerecord.Record; import com.jfinal.plugin.activerecord.Record;
import java.text.DateFormat;
import java.util.Date; import java.util.Date;
import java.util.HashMap;
import java.util.Map; import java.util.Map;
@Clear(AuthInterceptor.class) @Clear(AuthInterceptor.class)
@ -46,6 +45,9 @@ public class CacheController extends Controller {
map.set(entry.getKey().toString(), out); map.set(entry.getKey().toString(), out);
} }
map.set("task_queue_size", SyncTaskService.me.getTaskQueueSize());
renderJson(map); renderJson(map);
} }
} }

View File

@ -3,6 +3,7 @@ package com.cowr.service.ssjygl.jobs;
import com.cowr.common.Const; import com.cowr.common.Const;
import com.cowr.service.ssjygl.main.Config; import com.cowr.service.ssjygl.main.Config;
import com.cowr.service.ssjygl.main.SvrCacheData; import com.cowr.service.ssjygl.main.SvrCacheData;
import com.cowr.service.ssjygl.synctask.SyncTaskService;
import com.jfinal.kit.HttpKit; import com.jfinal.kit.HttpKit;
import com.jfinal.kit.StrKit; import com.jfinal.kit.StrKit;
import com.jfinal.log.Log; import com.jfinal.log.Log;
@ -61,6 +62,10 @@ public class CheckExceptionDataJob implements Job {
log.error("没有找到在线砂站信息。"); log.error("没有找到在线砂站信息。");
} }
if(SyncTaskService.me.getTaskQueueSize() > 10){
content += "task queue 还有 " + SyncTaskService.me.getTaskQueueSize() + " 条数据等待处理。";
}
if (StrKit.notBlank(content) && content.length() > 1) { if (StrKit.notBlank(content) && content.length() > 1) {
log.debug(content); log.debug(content);

View File

@ -89,7 +89,7 @@ public class Config extends JFinalConfig {
public static Prop dbprop = PropKit.use(ENV + "/db.properties", "UTF-8"); public static Prop dbprop = PropKit.use(ENV + "/db.properties", "UTF-8");
private WallFilter wallFilter; private WallFilter wallFilter;
public static NettyServer nettyServer = null; public static NettyServer nettyServer = null;
private static boolean server_run = true; public static boolean server_run = true;
private class ServerThread extends Thread { private class ServerThread extends Thread {
@Override @Override
@ -297,6 +297,8 @@ public class Config extends JFinalConfig {
} catch (Exception e) { } catch (Exception e) {
log.error(e.getMessage(), e); log.error(e.getMessage(), e);
} }
SyncTaskService.me.start();
} catch (Exception e) { } catch (Exception e) {
log.error(e.getMessage(), e); log.error(e.getMessage(), e);
} }

View File

@ -213,17 +213,18 @@ public class NettyServer {
} }
} else if (Enums.MsgTarget.SYNCTASK.name().equals(target)) { } else if (Enums.MsgTarget.SYNCTASK.name().equals(target)) {
JSONObject data = json.getJSONObject("data"); JSONObject data = json.getJSONObject("data");
boolean ret = SyncTaskService.me.recv(data, map.get(ctx.channel())); SyncTaskService.me.addTaskQueue(data);
// boolean ret = SyncTaskService.me.recv(data, map.get(ctx.channel()));
// 接收成功后返回id //
if (ret) { // // 接收成功后返回id
sendMsg(ctx, new JSONObject() // if (ret) {
.fluentPut("target", Enums.MsgTarget.SYNCRECV) // sendMsg(ctx, new JSONObject()
.fluentPut("id", data.get("id")) // .fluentPut("target", Enums.MsgTarget.SYNCRECV)
.toJSONString()); // .fluentPut("id", data.get("id"))
} else { // .toJSONString());
log.debug("数据接收后,入库失败:", msg); // } else {
} // log.debug("数据接收后,入库失败:", msg);
// }
} else if (Enums.MsgTarget.SYNCRECV.name().equals(target)) { } else if (Enums.MsgTarget.SYNCRECV.name().equals(target)) {
SyncTaskService.me.syncComplete(json.getString("id")); SyncTaskService.me.syncComplete(json.getString("id"));
} else if (Enums.MsgTarget.SYNCFAIL.name().equals(target)) { } else if (Enums.MsgTarget.SYNCFAIL.name().equals(target)) {

View File

@ -16,6 +16,8 @@ import com.jfinal.plugin.activerecord.Record;
import java.math.BigDecimal; import java.math.BigDecimal;
import java.sql.SQLException; import java.sql.SQLException;
import java.util.*; import java.util.*;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
/** /**
* Generated by COWR Sat Apr 04 19:57:12 CST 2020 * Generated by COWR Sat Apr 04 19:57:12 CST 2020
@ -26,12 +28,56 @@ import java.util.*;
public class SyncTaskService { public class SyncTaskService {
private static Log log = Log.getLog(SyncTaskService.class); private static Log log = Log.getLog(SyncTaskService.class);
public static final SyncTaskService me = new SyncTaskService(); public static final SyncTaskService me = new SyncTaskService();
private BlockingQueue<JSONObject> taskQueue = new LinkedBlockingQueue<>(); // 砂站推送到服务端的数据队列
// 是否开启 // 是否开启
public boolean isEnable() { public boolean isEnable() {
return CacheData.service_enable; return CacheData.service_enable;
} }
public void addTaskQueue(JSONObject obj) {
this.taskQueue.add(obj);
}
public int getTaskQueueSize(){
return this.taskQueue.size();
}
public void start() {
new TaskQueueThread().start();
}
private class TaskQueueThread extends Thread {
@Override
public void run() {
try {
while (Config.server_run) {
try {
JSONObject data = taskQueue.take();
int supermarket_id = data.getInteger("supermarket_id");
boolean ret = recv(data, supermarket_id);
// 接收成功后返回id
if (ret) {
Config.nettyServer.send(supermarket_id, new JSONObject()
.fluentPut("target", Enums.MsgTarget.SYNCRECV)
.fluentPut("id", data.get("id"))
.toJSONString());
} else {
log.debug("数据接收后,入库失败:", data);
}
Thread.sleep(100);
} catch (Exception e) {
log.error(e.getMessage(), e);
}
}
} catch (Exception e) {
log.error(e.getMessage(), e);
}
}
}
/** /**
* *
* *