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

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


The following commit(s) were added to refs/heads/master by this push:
     new f98b89e4a4 [INLONG-8396][Manager] Support for querying audit data with 
average delay (#8397)
f98b89e4a4 is described below

commit f98b89e4a4aa80b487d419e587d838cfdc29c7c2
Author: fuweng11 <[email protected]>
AuthorDate: Mon Jul 3 19:44:03 2023 +0800

    [INLONG-8396][Manager] Support for querying audit data with average delay 
(#8397)
---
 .../main/resources/mappers/AuditEntityMapper.xml   |  2 +-
 .../inlong/manager/pojo/audit/AuditInfo.java       |  7 ++-
 .../inlong/manager/pojo/audit/AuditRequest.java    |  8 ++--
 .../service/core/impl/AuditServiceImpl.java        | 52 +++++++++++++---------
 .../service/core/impl/AuditServiceTest.java        | 10 +++--
 .../manager/web/controller/AuditController.java    |  6 +--
 6 files changed, 52 insertions(+), 33 deletions(-)

diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/AuditEntityMapper.xml 
b/inlong-manager/manager-dao/src/main/resources/mappers/AuditEntityMapper.xml
index 1bc5d9ce82..ba8ffbf21b 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/AuditEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/AuditEntityMapper.xml
@@ -41,7 +41,7 @@
     </resultMap>
 
     <select id="sumByLogTs" resultMap="SumByLogTsResultMap">
-        select date_format(log_ts, #{format, jdbcType=VARCHAR}) as log_ts, 
sum(`count`) as total
+        select date_format(log_ts, #{format, jdbcType=VARCHAR}) as log_ts, 
sum(`count`) as total, sum(`delay`) as total_delay
         from apache_inlong_audit.audit_data
         where inlong_group_id = #{groupId,jdbcType=VARCHAR}
           and inlong_stream_id = #{streamId,jdbcType=VARCHAR}
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditInfo.java
index 6207a4b163..7a13088885 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditInfo.java
@@ -29,13 +29,16 @@ public class AuditInfo {
     @ApiModelProperty(value = "Audit log timestamp")
     private String logTs;
     @ApiModelProperty(value = "Audit count")
-    private Long count;
+    private long count;
+    @ApiModelProperty(value = "Audit delay")
+    private long delay;
 
     public AuditInfo() {
     }
 
-    public AuditInfo(String logTs, Long count) {
+    public AuditInfo(String logTs, long count, long delay) {
         this.logTs = logTs;
         this.count = count;
+        this.delay = delay;
     }
 }
\ No newline at end of file
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditRequest.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditRequest.java
index d9981dad7c..4023d40dc0 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditRequest.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/audit/AuditRequest.java
@@ -50,9 +50,11 @@ public class AuditRequest {
     @ApiModelProperty(value = "sink id")
     private Integer sinkId;
 
-    @ApiModelProperty(value = "query date, format by 'yyyy-MM-dd'", required = 
true, example = "2022-01-01")
-    @NotBlank(message = "dt not be blank")
-    private String dt;
+    @ApiModelProperty(value = "query start date, format by 'yyyy-MM-dd'", 
required = true, example = "2022-01-01")
+    private String startDate;
+
+    @ApiModelProperty(value = "query end date, format by 'yyyy-MM-dd'", 
required = true, example = "2022-01-01")
+    private String endDate;
 
     /**
      * Time statics dim such as MINUTE, HOUR, DAY
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 df18ea475b..df2cbc6748 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
@@ -73,6 +73,7 @@ import java.sql.SQLException;
 import java.text.SimpleDateFormat;
 import java.util.ArrayList;
 import java.util.Date;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.LinkedList;
 import java.util.List;
@@ -95,6 +96,9 @@ public class AuditServiceImpl implements AuditService {
     private static final String SECOND_FORMAT = "yyyy-MM-dd HH:mm:ss";
     private static final String HOUR_FORMAT = "yyyy-MM-dd HH";
     private static final String DAY_FORMAT = "yyyy-MM-dd";
+    private static final DateTimeFormatter SECOND_DATE_FORMATTER = 
DateTimeFormat.forPattern(SECOND_FORMAT);
+    private static final DateTimeFormatter HOUR_DATE_FORMATTER = 
DateTimeFormat.forPattern(HOUR_FORMAT);
+    private static final DateTimeFormatter DAY_DATE_FORMATTER = 
DateTimeFormat.forPattern(DAY_FORMAT);
 
     // key: type of audit base item, value: entity of audit base item
     private final Map<String, AuditBaseEntity> auditSentItemMap = new 
ConcurrentHashMap<>();
@@ -204,21 +208,21 @@ public class AuditServiceImpl implements AuditService {
             if (AuditQuerySource.MYSQL == querySource) {
                 String format = "%Y-%m-%d %H:%i:00";
                 // Support min agg at now
-                DateTimeFormatter forPattern = 
DateTimeFormat.forPattern("yyyy-MM-dd");
-                DateTime dtDate = forPattern.parseDateTime(request.getDt());
-                String eDate = dtDate.plusDays(1).toString(forPattern);
+                DateTime endDate = 
DAY_DATE_FORMATTER.parseDateTime(request.getEndDate());
+                String endDateStr = 
endDate.plusDays(1).toString(DAY_DATE_FORMATTER);
                 List<Map<String, Object>> sumList = 
auditEntityMapper.sumByLogTs(
-                        groupId, streamId, auditId, request.getDt(), eDate, 
format);
+                        groupId, streamId, auditId, request.getStartDate(), 
endDateStr, format);
                 List<AuditInfo> auditSet = sumList.stream().map(s -> {
                     AuditInfo vo = new AuditInfo();
                     vo.setLogTs((String) s.get("logTs"));
                     vo.setCount(((BigDecimal) s.get("total")).longValue());
+                    vo.setCount(((BigDecimal) 
s.get("total_delay")).longValue());
                     return vo;
                 }).collect(Collectors.toList());
                 result.add(new AuditVO(auditId, auditSet,
                         auditId.equals(getAuditId(sinkNodeType, true)) ? 
sinkNodeType : null));
             } else if (AuditQuerySource.ELASTICSEARCH == querySource) {
-                String index = String.format("%s_%s", 
request.getDt().replaceAll("-", ""), auditId);
+                String index = String.format("%s_%s", 
request.getStartDate().replaceAll("-", ""), auditId);
                 if (!elasticsearchApi.indexExists(index)) {
                     LOGGER.warn("elasticsearch index={} not exists", index);
                     continue;
@@ -232,6 +236,7 @@ public class AuditServiceImpl implements AuditService {
                             AuditInfo vo = new AuditInfo();
                             vo.setLogTs(bucket.getKeyAsString());
                             vo.setCount((long) ((ParsedSum) 
bucket.getAggregations().asList().get(0)).getValue());
+                            vo.setDelay((long) ((ParsedSum) 
bucket.getAggregations().asList().get(1)).getValue());
                             return vo;
                         }).collect(Collectors.toList());
                         result.add(new AuditVO(auditId, auditSet,
@@ -240,14 +245,15 @@ public class AuditServiceImpl implements AuditService {
                 }
             } else if (AuditQuerySource.CLICKHOUSE == querySource) {
                 try (Connection connection = 
ClickHouseConfig.getCkConnection();
-                        PreparedStatement statement =
-                                getAuditCkStatement(connection, groupId, 
streamId, auditId, request.getDt());
+                        PreparedStatement statement = 
getAuditCkStatement(connection, groupId, streamId, auditId,
+                                request.getStartDate(), request.getEndDate());
                         ResultSet resultSet = statement.executeQuery()) {
                     List<AuditInfo> auditSet = new ArrayList<>();
                     while (resultSet.next()) {
                         AuditInfo vo = new AuditInfo();
                         vo.setLogTs(resultSet.getString("log_ts"));
                         vo.setCount(resultSet.getLong("total"));
+                        vo.setDelay(resultSet.getLong("total_delay"));
                         auditSet.add(vo);
                     }
                     result.add(new AuditVO(auditId, auditSet,
@@ -295,7 +301,8 @@ public class AuditServiceImpl implements AuditService {
      */
     private SearchRequest toAuditSearchRequest(String index, String groupId, 
String streamId) {
         TermsAggregationBuilder builder = 
AggregationBuilders.terms("log_ts").field("log_ts")
-                
.size(Integer.MAX_VALUE).subAggregation(AggregationBuilders.sum("count").field("count"));
+                
.size(Integer.MAX_VALUE).subAggregation(AggregationBuilders.sum("count").field("count"))
+                
.subAggregation(AggregationBuilders.sum("delay").field("delay"));
         BoolQueryBuilder filterBuilder = new BoolQueryBuilder();
         filterBuilder.must(termQuery("inlong_group_id", groupId));
         filterBuilder.must(termQuery("inlong_stream_id", streamId));
@@ -311,22 +318,20 @@ public class AuditServiceImpl implements AuditService {
     /**
      * Get clickhouse Statement
      *
-     * @param connection The ClickHouse connection
      * @param groupId The groupId of inlong
      * @param streamId The streamId of inlong
      * @param auditId The auditId of request
-     * @param dt The datetime of request
+     * @param startDate The start datetime of request
+     * @param endDate The en datetime of request
      * @return The clickhouse Statement
      */
     private PreparedStatement getAuditCkStatement(Connection connection, 
String groupId, String streamId,
-            String auditId, String dt) throws SQLException {
-        DateTimeFormatter formatter = DateTimeFormat.forPattern(DAY_FORMAT);
-        DateTime date = formatter.parseDateTime(dt);
-        String startDate = date.toString(SECOND_FORMAT);
-        String endDate = date.plusDays(1).toString(SECOND_FORMAT);
+            String auditId, String startDate, String endDate) throws 
SQLException {
+        String start = 
DAY_DATE_FORMATTER.parseDateTime(startDate).toString(SECOND_FORMAT);
+        String end = 
DAY_DATE_FORMATTER.parseDateTime(endDate).plusDays(1).toString(SECOND_FORMAT);
 
         String sql = new SQL()
-                .SELECT("log_ts", "sum(count) as total")
+                .SELECT("log_ts", "sum(count) as total", "sum(delay) as 
total_delay")
                 .FROM("audit_data")
                 .WHERE("inlong_group_id = ?")
                 .WHERE("inlong_stream_id = ?")
@@ -341,8 +346,8 @@ public class AuditServiceImpl implements AuditService {
         statement.setString(1, groupId);
         statement.setString(2, streamId);
         statement.setString(3, auditId);
-        statement.setString(4, startDate);
-        statement.setString(5, endDate);
+        statement.setString(4, start);
+        statement.setString(5, end);
         return statement;
     }
 
@@ -359,7 +364,7 @@ public class AuditServiceImpl implements AuditService {
                 result = doAggregate(auditVOList, DAY_FORMAT);
                 break;
             default:
-                result = auditVOList;
+                result = doAggregate(auditVOList, SECOND_FORMAT);
                 break;
         }
         return result;
@@ -372,7 +377,8 @@ public class AuditServiceImpl implements AuditService {
         List<AuditVO> result = new ArrayList<>();
         for (AuditVO auditVO : auditVOList) {
             AuditVO statInfo = new AuditVO();
-            ConcurrentHashMap<String, AtomicLong> countMap = new 
ConcurrentHashMap<>();
+            HashMap<String, AtomicLong> countMap = new HashMap<>();
+            HashMap<String, AtomicLong> delayMap = new HashMap<>();
             statInfo.setAuditId(auditVO.getAuditId());
             statInfo.setNodeType(auditVO.getNodeType());
             for (AuditInfo auditInfo : auditVO.getAuditSet()) {
@@ -383,14 +389,20 @@ public class AuditServiceImpl implements AuditService {
                 if (countMap.get(statKey) == null) {
                     countMap.put(statKey, new AtomicLong(0));
                 }
+                if (delayMap.get(statKey) == null) {
+                    delayMap.put(statKey, new AtomicLong(0));
+                }
                 countMap.get(statKey).addAndGet(auditInfo.getCount());
+                delayMap.get(statKey).addAndGet(auditInfo.getDelay());
             }
 
             List<AuditInfo> auditInfoList = new LinkedList<>();
             for (Map.Entry<String, AtomicLong> entry : countMap.entrySet()) {
                 AuditInfo auditInfoStat = new AuditInfo();
                 auditInfoStat.setLogTs(entry.getKey());
+                long count = entry.getValue().get();
                 auditInfoStat.setCount(entry.getValue().get());
+                auditInfoStat.setDelay(count == 0 ? 0 : 
delayMap.get(entry.getKey()).get() / count);
                 auditInfoList.add(auditInfoStat);
             }
             statInfo.setAuditSet(auditInfoList);
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AuditServiceTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AuditServiceTest.java
index 12fed2ff36..010a1bd38a 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AuditServiceTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AuditServiceTest.java
@@ -46,12 +46,13 @@ class AuditServiceTest extends ServiceBaseTest {
         request.setAuditIds(Arrays.asList("3", "4"));
         request.setInlongGroupId("g1");
         request.setInlongStreamId("s1");
-        request.setDt("2022-01-01");
+        request.setStartDate("2022-01-01");
+        request.setEndDate("2022-01-01");
         List<AuditVO> result = new ArrayList<>();
         AuditVO auditVO = new AuditVO();
         auditVO.setAuditId("3");
-        auditVO.setAuditSet(Arrays.asList(new AuditInfo("2022-01-01 00:00:00", 
123L),
-                new AuditInfo("2022-01-01 00:01:00", 124L)));
+        auditVO.setAuditSet(Arrays.asList(new AuditInfo("2022-01-01 00:00:00", 
123L, 12L),
+                new AuditInfo("2022-01-01 00:01:00", 124L, 12L)));
         result.add(auditVO);
         Assertions.assertNotNull(result);
         // close real test for testQueryFromMySQL due to date_format function 
not support in h2
@@ -70,7 +71,8 @@ class AuditServiceTest extends ServiceBaseTest {
         request.setAuditIds(Arrays.asList("3", "4"));
         request.setInlongGroupId("g1");
         request.setInlongStreamId("s1");
-        request.setDt("2022-01-01");
+        request.setStartDate("2022-01-01");
+        request.setEndDate("2022-01-01");
         Assertions.assertNotNull(auditService.listByCondition(request));
     }
 }
diff --git 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/AuditController.java
 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/AuditController.java
index 9d78005eeb..1e76796cf2 100644
--- 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/AuditController.java
+++ 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/AuditController.java
@@ -26,8 +26,8 @@ import io.swagger.annotations.Api;
 import io.swagger.annotations.ApiOperation;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.context.annotation.Lazy;
-import org.springframework.web.bind.annotation.GetMapping;
 import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RestController;
 
@@ -47,9 +47,9 @@ public class AuditController {
     @Autowired
     private AuditService auditService;
 
-    @GetMapping(value = "/audit/list")
+    @PostMapping(value = "/audit/list")
     @ApiOperation(value = "Query audit list according to conditions")
-    public Response<List<AuditVO>> listByCondition(@Valid AuditRequest 
request) throws Exception {
+    public Response<List<AuditVO>> listByCondition(@Valid @RequestBody 
AuditRequest request) throws Exception {
         return Response.success(auditService.listByCondition(request));
     }
 

Reply via email to