From 92aa25940cb4ddef6c2ad7934f5bd6aed365a2bf Mon Sep 17 00:00:00 2001 From: aoshiguchen <1052045476@qq.com> Date: Thu, 27 Oct 2022 21:39:08 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=8C=E5=96=84=E3=80=90=E6=B5=81=E9=87=8F?= =?UTF-8?q?=E7=BB=9F=E8=AE=A1=E6=8A=A5=E8=A1=A8=20-=20=E5=88=86=E9=92=9F?= =?UTF-8?q?=E7=BA=A7=E5=88=AB=E3=80=91=E5=A4=84=E7=90=86=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../server/dal/FlowReportMinuteMapper.java | 14 ++++++++++++++ .../server/job/FlowReportForMinuteJob.java | 19 +++++++++++++++---- .../proxy/core/BytesMetricsHandler.java | 2 +- .../proxy/core/VisitorChannelHandler.java | 2 +- .../handler/ProxyMessageTransferHandler.java | 2 +- 5 files changed, 32 insertions(+), 7 deletions(-) diff --git a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/dal/FlowReportMinuteMapper.java b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/dal/FlowReportMinuteMapper.java index 35762541..1f02dccd 100644 --- a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/dal/FlowReportMinuteMapper.java +++ b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/dal/FlowReportMinuteMapper.java @@ -22,11 +22,17 @@ package fun.asgc.neutrino.proxy.server.dal; import fun.asgc.neutrino.core.annotation.Component; +import fun.asgc.neutrino.core.annotation.Param; import fun.asgc.neutrino.core.aop.Intercept; import fun.asgc.neutrino.core.db.annotation.Insert; +import fun.asgc.neutrino.core.db.annotation.ResultType; +import fun.asgc.neutrino.core.db.annotation.Select; import fun.asgc.neutrino.core.db.mapper.SqlMapper; import fun.asgc.neutrino.proxy.server.dal.entity.FlowReportMinuteDO; +import java.util.List; +import java.util.Set; + /** * @author: aoshiguchen * @date: 2022/10/24 @@ -34,5 +40,13 @@ import fun.asgc.neutrino.proxy.server.dal.entity.FlowReportMinuteDO; @Intercept(ignoreGlobal = true) @Component public interface FlowReportMinuteMapper extends SqlMapper { + @Select("select * from flow_report_minute where license_id = :licenseId and date = :date") + FlowReportMinuteDO findOne(@Param("licenseId") Integer licenseId, @Param("date") String date); + @ResultType(FlowReportMinuteDO.class) + @Select("select * from flow_report_minute where license_id in (:licenseIds) and date = :date") + List findList(@Param("licenseIds") Set licenseIds, @Param("date") String date); + + @Insert("insert into flow_report_minute(`user_id`,`license_id`,`write_bytes`,`read_bytes`,`date`,`create_time`) values(:userId,:licenseId,:writeBytes,:readBytes,:date,:createTime)") + void add(FlowReportMinuteDO flowReportMinuteDO); } diff --git a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/job/FlowReportForMinuteJob.java b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/job/FlowReportForMinuteJob.java index fef65ec4..47922b35 100644 --- a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/job/FlowReportForMinuteJob.java +++ b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/job/FlowReportForMinuteJob.java @@ -35,8 +35,9 @@ import fun.asgc.neutrino.proxy.server.dal.entity.LicenseDO; import fun.asgc.neutrino.proxy.server.service.FlowReportService; import lombok.extern.slf4j.Slf4j; -import java.util.Date; -import java.util.List; +import java.util.*; +import java.util.function.Function; +import java.util.stream.Collectors; /** * 流量统计报表 - 分钟级别 @@ -62,8 +63,18 @@ public class FlowReportForMinuteJob implements IJobHandler { if (CollectionUtil.isEmpty(list)) { return; } + Set licenseIds = list.stream().map(LicenseDO::getId).collect(Collectors.toSet()); Date now = new Date(); + String date = DateUtil.format(DateUtil.addDate(now, Calendar.MINUTE, -1), "yyyy-MM-dd HH:mm"); + List oldList = flowReportMinuteMapper.findList(licenseIds, date); + Map oldMap = CollectionUtil.isEmpty(oldList) ? new HashMap<>() : + oldList.stream().collect(Collectors.toMap(FlowReportMinuteDO::getLicenseId, Function.identity(), (a,b) -> a)); + for (LicenseDO item : list) { + // 避免job重复执行导致数据重复 + if (oldMap.containsKey(item.getId())) { + continue; + } Integer writeBytes = flowReportService.getAndResetWriteByte(item.getId()); Integer readBytes = flowReportService.getAndResetReadByte(item.getId()); if (writeBytes == 0 && readBytes == 0) { @@ -74,9 +85,9 @@ public class FlowReportForMinuteJob implements IJobHandler { flowReportMinuteDO.setLicenseId(item.getId()); flowReportMinuteDO.setWriteBytes(writeBytes); flowReportMinuteDO.setReadBytes(readBytes); - flowReportMinuteDO.setDate(DateUtil.format(now, "yyyy-MM-dd HH:mm")); + flowReportMinuteDO.setDate(date); flowReportMinuteDO.setCreateTime(now); - // TODO insert + flowReportMinuteMapper.add(flowReportMinuteDO); } } } diff --git a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/BytesMetricsHandler.java b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/BytesMetricsHandler.java index 4182d9ca..d157376c 100644 --- a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/BytesMetricsHandler.java +++ b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/BytesMetricsHandler.java @@ -43,7 +43,7 @@ public class BytesMetricsHandler extends ChannelDuplexHandler { MetricsCollector metricsCollector = MetricsCollector.getCollector(sa.getPort()); metricsCollector.incrementReadBytes(((ByteBuf) msg).readableBytes()); metricsCollector.incrementReadMsgs(1); - System.out.println("字节数:" + metricsCollector.getMetrics().getReadBytes()); +// System.out.println("字节数:" + metricsCollector.getMetrics().getReadBytes()); ctx.fireChannelRead(msg); } diff --git a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/VisitorChannelHandler.java b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/VisitorChannelHandler.java index de829dca..481f2326 100755 --- a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/VisitorChannelHandler.java +++ b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/core/VisitorChannelHandler.java @@ -71,7 +71,7 @@ public class VisitorChannelHandler extends SimpleChannelInboundHandler // 增加流量计数 VisitorChannelAttachInfo visitorChannelAttachInfo = ProxyUtil.getAttachInfo(visitorChannel); - BeanManager.getBean(FlowReportService.class).addWriteByte(visitorChannelAttachInfo.getLicenseId(), buf.readableBytes()); + BeanManager.getBean(FlowReportService.class).addWriteByte(visitorChannelAttachInfo.getLicenseId(), bytes.length); } } diff --git a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/handler/ProxyMessageTransferHandler.java b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/handler/ProxyMessageTransferHandler.java index 367e07fd..65e69554 100644 --- a/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/handler/ProxyMessageTransferHandler.java +++ b/neutrino-proxy-server/src/main/java/fun/asgc/neutrino/proxy/server/proxy/handler/ProxyMessageTransferHandler.java @@ -57,7 +57,7 @@ public class ProxyMessageTransferHandler implements ProxyMessageHandler { // 增加流量计数 VisitorChannelAttachInfo visitorChannelAttachInfo = ProxyUtil.getAttachInfo(visitorChannel); - BeanManager.getBean(FlowReportService.class).addReadByte(visitorChannelAttachInfo.getLicenseId(), buf.readableBytes()); + BeanManager.getBean(FlowReportService.class).addReadByte(visitorChannelAttachInfo.getLicenseId(), proxyMessage.getData().length); } }