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