133 lines
4.6 KiB
Go
133 lines
4.6 KiB
Go
package repository
|
||
|
||
import (
|
||
"fmt"
|
||
"strings"
|
||
|
||
"be.ems/src/framework/datasource"
|
||
"be.ems/src/framework/logger"
|
||
"be.ems/src/modules/network_data/model"
|
||
)
|
||
|
||
// 实例化数据层 PerfKPIImpl 结构体
|
||
var NewPerfKPIImpl = &PerfKPIImpl{}
|
||
|
||
// PerfKPIImpl 性能统计 数据层处理
|
||
type PerfKPIImpl struct{}
|
||
|
||
// SelectGoldKPI 通过网元指标数据信息
|
||
func (r *PerfKPIImpl) SelectGoldKPI(query model.GoldKPIQuery, kpiIds []string) []map[string]any {
|
||
// 查询条件拼接
|
||
var conditions []string
|
||
var params []any
|
||
var tableName string = "kpi_report_"
|
||
if query.RmUID != "" {
|
||
conditions = append(conditions, "gk.rm_uid = ?")
|
||
params = append(params, query.RmUID)
|
||
}
|
||
if query.NeType != "" {
|
||
//conditions = append(conditions, "gk.ne_type = ?")
|
||
// params = append(params, query.NeType)
|
||
tableName += strings.ToLower(query.NeType)
|
||
}
|
||
if query.StartTime != "" {
|
||
conditions = append(conditions, "gk.created_at >= ?")
|
||
params = append(params, query.StartTime)
|
||
}
|
||
if query.EndTime != "" {
|
||
conditions = append(conditions, "gk.created_at <= ?")
|
||
params = append(params, query.EndTime)
|
||
}
|
||
|
||
// 构建查询条件语句
|
||
whereSql := ""
|
||
if len(conditions) > 0 {
|
||
whereSql += " where " + strings.Join(conditions, " and ")
|
||
}
|
||
|
||
// 查询字段列
|
||
var fields = []string{
|
||
// fmt.Sprintf("FROM_UNIXTIME(FLOOR(gk.created_at / (%d * 1000)) * %d) AS timeGroup", query.Interval, query.Interval),
|
||
fmt.Sprintf("CONCAT(FLOOR(gk.created_at / (%d * 1000)) * (%d * 1000)) AS timeGroup", query.Interval, query.Interval), // 时间戳毫秒
|
||
"min(CASE WHEN gk.index != '' THEN gk.index ELSE 0 END) AS startIndex",
|
||
"min(CASE WHEN gk.ne_type != '' THEN gk.ne_type ELSE 0 END) AS neType",
|
||
"min(CASE WHEN gk.ne_name != '' THEN gk.ne_name ELSE 0 END) AS neName",
|
||
"min(CASE WHEN gk.rm_uid != '' THEN gk.rm_uid ELSE 0 END) AS rmUID",
|
||
}
|
||
for i, kid := range kpiIds {
|
||
// 特殊字段,只取最后一次收到的非0值
|
||
if kid == "AMF.01" || kid == "UDM.01" || kid == "UDM.02" || kid == "UDM.03" || kid == "SMF.01" {
|
||
str := fmt.Sprintf("IFNULL(SUBSTRING_INDEX(GROUP_CONCAT( CASE WHEN JSON_EXTRACT(gk.kpi_values, '$[%d].kpi_id') = '%s' THEN JSON_EXTRACT(gk.kpi_values, '$[%d].value') END ), ',', 1), 0) AS '%s'", i, kid, i, kid)
|
||
fields = append(fields, str)
|
||
} else {
|
||
str := fmt.Sprintf("sum(CASE WHEN JSON_EXTRACT(gk.kpi_values, '$[%d].kpi_id') = '%s' THEN JSON_EXTRACT(gk.kpi_values, '$[%d].value') ELSE 0 END) AS '%s'", i, kid, i, kid)
|
||
fields = append(fields, str)
|
||
}
|
||
}
|
||
fieldsSql := strings.Join(fields, ",")
|
||
|
||
// 查询数据
|
||
if query.SortField == "" {
|
||
query.SortField = "timeGroup"
|
||
}
|
||
if query.SortOrder == "" {
|
||
query.SortOrder = "desc"
|
||
}
|
||
orderSql := fmt.Sprintf(" order by %s %s", query.SortField, query.SortOrder)
|
||
querySql := fmt.Sprintf("SELECT %s FROM %s gk %s GROUP BY timeGroup %s", fieldsSql, tableName, whereSql, orderSql)
|
||
results, err := datasource.RawDB("", querySql, params)
|
||
if err != nil {
|
||
logger.Errorf("query err => %v", err)
|
||
}
|
||
return results
|
||
}
|
||
|
||
// SelectGoldKPITitle 网元对应的指标名称
|
||
func (r *PerfKPIImpl) SelectGoldKPITitle(neType string) []model.GoldKPITitle {
|
||
result := []model.GoldKPITitle{}
|
||
tx := datasource.DefaultDB().Table("kpi_title").Where("ne_type = ?", neType).Find(&result)
|
||
if err := tx.Error; err != nil {
|
||
logger.Errorf("Find err => %v", err)
|
||
}
|
||
return result
|
||
}
|
||
|
||
// SelectUPFTotalFlow 查询UPF总流量 N3上行 N6下行
|
||
func (r *PerfKPIImpl) SelectUPFTotalFlow(neType, rmUID, startDate, endDate string) map[string]any {
|
||
// 查询条件拼接
|
||
var conditions []string
|
||
var params []any
|
||
if neType != "" {
|
||
conditions = append(conditions, "kupf.ne_type = ?")
|
||
params = append(params, neType)
|
||
}
|
||
if rmUID != "" {
|
||
conditions = append(conditions, "kupf.rm_uid = ?")
|
||
params = append(params, rmUID)
|
||
}
|
||
if startDate != "" {
|
||
conditions = append(conditions, "kupf.created_at >= ?")
|
||
params = append(params, startDate)
|
||
}
|
||
if endDate != "" {
|
||
conditions = append(conditions, "kupf.created_at <= ?")
|
||
params = append(params, endDate)
|
||
}
|
||
// 构建查询条件语句
|
||
whereSql := ""
|
||
if len(conditions) > 0 {
|
||
whereSql += " where " + strings.Join(conditions, " and ")
|
||
}
|
||
|
||
// 查询数据
|
||
querySql := `SELECT
|
||
sum( CASE WHEN JSON_EXTRACT(kupf.kpi_values, '$[2].kpi_id') = 'UPF.03' THEN JSON_EXTRACT(kupf.kpi_values, '$[2].value') ELSE 0 END ) AS 'up',
|
||
sum( CASE WHEN JSON_EXTRACT(kupf.kpi_values, '$[5].kpi_id') = 'UPF.06' THEN JSON_EXTRACT(kupf.kpi_values, '$[5].value') ELSE 0 END ) AS 'down'
|
||
FROM kpi_report_upf kupf`
|
||
results, err := datasource.RawDB("", querySql+whereSql, params)
|
||
if err != nil {
|
||
logger.Errorf("query err => %v", err)
|
||
}
|
||
return results[0]
|
||
}
|