This is an automated email from the ASF dual-hosted git repository.

healchow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new d9d7b1455 [INLONG-4656][Manager] Support query metrics by day and hour 
for Audit (#4658)
d9d7b1455 is described below

commit d9d7b14550f087d42fe00f2b8ad03784191c8019
Author: doleyzi <[email protected]>
AuthorDate: Wed Jun 15 14:02:34 2022 +0800

    [INLONG-4656][Manager] Support query metrics by day and hour for Audit 
(#4658)
---
 .../service/core/impl/AuditServiceImpl.java        | 145 ++++++++++++++++-----
 1 file changed, 109 insertions(+), 36 deletions(-)

diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AuditServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AuditServiceImpl.java
index 41588ca9e..0247b0229 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AuditServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AuditServiceImpl.java
@@ -17,16 +17,9 @@
 
 package org.apache.inlong.manager.service.core.impl;
 
-import static org.elasticsearch.index.query.QueryBuilders.termQuery;
-
-import java.io.IOException;
-import java.math.BigDecimal;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
-import java.util.stream.Collectors;
 import org.apache.commons.collections.CollectionUtils;
 import org.apache.inlong.manager.common.enums.AuditQuerySource;
+import org.apache.inlong.manager.common.enums.TimeStaticsDim;
 import org.apache.inlong.manager.common.pojo.audit.AuditInfo;
 import org.apache.inlong.manager.common.pojo.audit.AuditRequest;
 import org.apache.inlong.manager.common.pojo.audit.AuditVO;
@@ -53,6 +46,20 @@ import 
org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.beans.factory.annotation.Value;
 import org.springframework.stereotype.Service;
 
+import java.io.IOException;
+import java.math.BigDecimal;
+import java.text.SimpleDateFormat;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.LinkedList;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.stream.Collectors;
+
+import static org.elasticsearch.index.query.QueryBuilders.termQuery;
+
 /**
  * Audit service layer implementation
  */
@@ -60,10 +67,11 @@ import org.springframework.stereotype.Service;
 public class AuditServiceImpl implements AuditService {
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(AuditServiceImpl.class);
+    private static final String HOUR_FORMAT = "yyyy-MM-dd HH";
+    private static final String DAY_FORMAT = "yyyy-MM-dd";
 
     @Value("${audit.query.source}")
     private String auditQuerySource = AuditQuerySource.MYSQL.name();
-
     @Autowired
     private AuditEntityMapper auditEntityMapper;
     @Autowired
@@ -73,6 +81,9 @@ public class AuditServiceImpl implements AuditService {
     public List<AuditVO> listByCondition(AuditRequest request) throws 
IOException {
         LOGGER.info("begin query audit list request={}", request);
         Preconditions.checkNotNull(request, "request is null");
+
+        String groupId = request.getInlongGroupId();
+        String streamId = request.getInlongStreamId();
         List<AuditVO> result = new ArrayList<>();
         AuditQuerySource querySource = 
AuditQuerySource.valueOf(auditQuerySource);
         for (String auditId : request.getAuditIds()) {
@@ -82,9 +93,8 @@ public class AuditServiceImpl implements AuditService {
                 DateTimeFormatter forPattern = 
DateTimeFormat.forPattern("yyyy-MM-dd");
                 DateTime dtDate = forPattern.parseDateTime(request.getDt());
                 String eDate = dtDate.plusDays(1).toString(forPattern);
-                List<Map<String, Object>> sumList = auditEntityMapper
-                        .sumByLogTs(request.getInlongGroupId(), 
request.getInlongStreamId(),
-                                auditId, request.getDt(), eDate, format);
+                List<Map<String, Object>> sumList = 
auditEntityMapper.sumByLogTs(
+                        groupId, streamId, auditId, request.getDt(), eDate, 
format);
                 List<AuditInfo> auditSet = sumList.stream().map(s -> {
                     AuditInfo vo = new AuditInfo();
                     vo.setLogTs((String) s.get("logTs"));
@@ -93,32 +103,29 @@ public class AuditServiceImpl implements AuditService {
                 }).collect(Collectors.toList());
                 result.add(new AuditVO(auditId, auditSet));
             } else if (AuditQuerySource.ELASTICSEARCH == querySource) {
-                String index = String.format("%s_%s",
-                        request.getDt().replaceAll("-", ""), auditId);
-                if (elasticsearchApi.indexExists(index)) {
-                    SearchResponse response = elasticsearchApi
-                            .search(toAuditSearchRequest(index, 
request.getInlongGroupId(),
-                                    request.getInlongStreamId()));
-                    final List<Aggregation> aggregations = 
response.getAggregations().asList();
-                    if (CollectionUtils.isNotEmpty(aggregations)) {
-                        ParsedTerms terms = (ParsedTerms) aggregations.get(0);
-                        if (CollectionUtils.isNotEmpty(terms.getBuckets())) {
-                            List<AuditInfo> auditSet = 
terms.getBuckets().stream().map(bucket -> {
-                                AuditInfo vo = new AuditInfo();
-                                vo.setLogTs(bucket.getKeyAsString());
-                                vo.setCount((long) ((ParsedSum) 
bucket.getAggregations().asList().get(0)).getValue());
-                                return vo;
-                            }).collect(Collectors.toList());
-                            result.add(new AuditVO(auditId, auditSet));
-                        }
+                String index = String.format("%s_%s", 
request.getDt().replaceAll("-", ""), auditId);
+                if (!elasticsearchApi.indexExists(index)) {
+                    LOGGER.warn("elasticsearch index={} not exists", index);
+                    continue;
+                }
+                SearchResponse response = 
elasticsearchApi.search(toAuditSearchRequest(index, groupId, streamId));
+                final List<Aggregation> aggregations = 
response.getAggregations().asList();
+                if (CollectionUtils.isNotEmpty(aggregations)) {
+                    ParsedTerms terms = (ParsedTerms) aggregations.get(0);
+                    if (CollectionUtils.isNotEmpty(terms.getBuckets())) {
+                        List<AuditInfo> auditSet = 
terms.getBuckets().stream().map(bucket -> {
+                            AuditInfo vo = new AuditInfo();
+                            vo.setLogTs(bucket.getKeyAsString());
+                            vo.setCount((long) ((ParsedSum) 
bucket.getAggregations().asList().get(0)).getValue());
+                            return vo;
+                        }).collect(Collectors.toList());
+                        result.add(new AuditVO(auditId, auditSet));
                     }
-                } else {
-                    LOGGER.warn("Elasticsearch index={} not exists", index);
                 }
             }
         }
         LOGGER.info("success to query audit list for request={}", request);
-        return result;
+        return aggregateByTimeDim(result, request.getTimeStaticsDim());
     }
 
     /**
@@ -130,13 +137,13 @@ public class AuditServiceImpl implements AuditService {
      * @return The search request of elasticsearch
      */
     private SearchRequest toAuditSearchRequest(String index, String groupId, 
String streamId) {
-        TermsAggregationBuilder aggrBuilder = 
AggregationBuilders.terms("log_ts").field("log_ts")
+        TermsAggregationBuilder builder = 
AggregationBuilders.terms("log_ts").field("log_ts")
                 
.size(Integer.MAX_VALUE).subAggregation(AggregationBuilders.sum("count").field("count"));
         BoolQueryBuilder filterBuilder = new BoolQueryBuilder();
         filterBuilder.must(termQuery("inlong_group_id", groupId));
         filterBuilder.must(termQuery("inlong_stream_id", streamId));
         SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
-        sourceBuilder.aggregation(aggrBuilder);
+        sourceBuilder.aggregation(builder);
         sourceBuilder.query(filterBuilder);
         sourceBuilder.from(0);
         sourceBuilder.size(0);
@@ -144,4 +151,70 @@ public class AuditServiceImpl implements AuditService {
         return new SearchRequest(new String[]{index}, sourceBuilder);
     }
 
-}
\ No newline at end of file
+    /**
+     * Aggregate by time dim
+     */
+    private List<AuditVO> aggregateByTimeDim(List<AuditVO> auditVOList, 
TimeStaticsDim timeStaticsDim) {
+        List<AuditVO> result;
+        switch (timeStaticsDim) {
+            case HOUR:
+                result = doAggregate(auditVOList, HOUR_FORMAT);
+                break;
+            case DAY:
+                result = doAggregate(auditVOList, DAY_FORMAT);
+                break;
+            default:
+                result = auditVOList;
+                break;
+        }
+        return result;
+    }
+
+    /**
+     * Execute the aggregate by the given time format
+     */
+    private List<AuditVO> doAggregate(List<AuditVO> auditVOList, String 
format) {
+        List<AuditVO> result = new ArrayList<>();
+        for (AuditVO auditVO : auditVOList) {
+            AuditVO statInfo = new AuditVO();
+            ConcurrentHashMap<String, AtomicLong> countMap = new 
ConcurrentHashMap<>();
+            statInfo.setAuditId(auditVO.getAuditId());
+            for (AuditInfo auditInfo : auditVO.getAuditSet()) {
+                String statKey = formatLogTime(auditInfo.getLogTs(), format);
+                if (statKey == null) {
+                    continue;
+                }
+                if (countMap.get(statKey) == null) {
+                    countMap.put(statKey, new AtomicLong(0));
+                }
+                countMap.get(statKey).addAndGet(auditInfo.getCount());
+            }
+
+            List<AuditInfo> auditInfoList = new LinkedList<>();
+            for (Map.Entry<String, AtomicLong> entry : countMap.entrySet()) {
+                AuditInfo auditInfoStat = new AuditInfo();
+                auditInfoStat.setLogTs(entry.getKey());
+                auditInfoStat.setCount(entry.getValue().get());
+                auditInfoList.add(auditInfoStat);
+            }
+            statInfo.setAuditSet(auditInfoList);
+            result.add(statInfo);
+        }
+        return result;
+    }
+
+    /**
+     * Format the log time
+     */
+    private String formatLogTime(String dateString, String format) {
+        String formatDateString = null;
+        try {
+            SimpleDateFormat formatter = new SimpleDateFormat(format);
+            Date date = formatter.parse(dateString);
+            formatDateString = formatter.format(date);
+        } catch (Exception e) {
+            LOGGER.error("format lot time exception", e);
+        }
+        return formatDateString;
+    }
+}

Reply via email to