引入多数据源及内存数据库

This commit is contained in:
zhouxiunai
2018-11-28 20:32:19 +08:00
parent 799544ff9c
commit f8d392d5af
23 changed files with 424 additions and 72 deletions
@@ -1,18 +1,12 @@
package com.gzzn.omms.msgexchangeapi.task;
import java.util.List;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.PageRequest;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import com.gzzn.omms.msgexchangeapi.dao.CminmsgDao;
import com.gzzn.omms.msgexchangeapi.entiy.Cminmsg;
import com.gzzn.omms.msgexchangeapi.exception.ExchangeServiceException;
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
import com.gzzn.omms.msgexchangeapi.service.IExchangeService;
import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
@@ -32,39 +26,14 @@ public class ExchangeTask {
IExchangeService exchangeService;
@Autowired
private CminmsgDao cminmsgDao;
private ICminmsgService cminmsgService;
@Scheduled(cron="0 0/1 * * * ?")
public void corn()
{
logger.info("定时任务启动....");
Long maxPerLong = 1000L; //每次最大发送多少条消息
Long lenTotal = cminmsgDao.getCminmsgsDateProcessedIsNotNullCount();
Long lenSend = 0L;
Integer pageIndex = 0;
Integer pageSize = 50;
while(lenSend < lenTotal && lenSend < maxPerLong)
{
List<Cminmsg> ls = cminmsgDao.findByCminmsgsDateProcessedIsNotNull(new PageRequest(pageIndex,pageSize));
for(Cminmsg node:ls)
{
String strJson = null;
try {
strJson = exchangeService.xmlToJson(node.getCminmsgsClobMsg());
if(!StringUtils.isEmpty(strJson))
{
kafkaservice.msgSend("topic2", strJson);
}
} catch (ExchangeServiceException e) {
logger.error(e.getMessage());
}
} // end for
lenSend += ls.size();
pageIndex++;
}//end while
//
logger.info("同步完成");
}