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;
+ }
+}