Files
qb/mqe/unaggregated.go
2026-06-19 07:25:46 +03:00

139 lines
4.3 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package mqe
import (
"fmt"
"sync"
"time"
)
// Measure - рассчитанное показание реальной метрики. Рассчитанное означает
// что учтен Factor. То есть значение уже можно показывать юзеру
// Для instant - одно значение value, для cumulative - 2 значения (value, total)
type Measure struct {
Time int64 `json:"t"`
Values []float64 `json:"v"`
}
type UnaggregatedFilter struct {
Since time.Time
Until time.Time
LastDayOfMonth int // например, конец месяца 25 число
FirstHourOfDay int // например день начинается в 6:00
InstantMetrics []int64
CumulativeMetrics []int64
}
// type MetricsDataResultMaps struct {
// AggregatedCumulative map[int64][]CalculatedAggregatedMeasure
// AggregatedByFuncInstant map[int64][]AggregatedByFuncInstantMeasure
// Cumulative map[int64][]CalculatedMeasure
// Instant map[int64][]CalculatedMeasure
// }
// SelectObjectDataInParallel - метод достает из БД данные всех метрик, которые связаны
// с объектом. Учитываются настройки объекта (LastDayOfMonth, FirstHourOfDay), а также
// фильтры выбранные пользователем на странице объекта (Since, Until, GroupBy).
// ВАЖНО!
// Запросы по каждой метрике отправляются в БД в отдельной горутине, что позволяет
// в разы ускорить получение данных по объекту.
func (s *MeasureQueryEngine) SelectUnaggregatedData(req UnaggregatedFilter) (resultMap map[int64][]Measure, err error) {
since := req.Since
until := req.Until
if req.FirstHourOfDay < 0 || req.FirstHourOfDay > 23 {
err = fix.Field(InvalidFirstHourOfDay, "firstHourOfDay")
return
}
if req.LastDayOfMonth < 0 || req.LastDayOfMonth > 28 {
err = fix.Field(InvalidLastDayOfMonth, "lastDayOfMonth")
return
}
// fmt.Printf(`SelectUnaggregatedData: {
// Since: %s
// Until: %s
// FirstHourOfDay: %d
// LastDayOfMonth: %d
// }
// `, since, until, req.FirstHourOfDay, req.LastDayOfMonth)
// FIX проверить как добавляются часы в DST часовых поясах
// Для req.FirstHourOfDay=4 since должен быть 4:00:00, until должен быть
// 3:59:59 следующего дня
if req.FirstHourOfDay > 0 {
since = since.Add(time.Duration(req.FirstHourOfDay) * time.Hour)
until = until.Add(time.Duration(req.FirstHourOfDay) * time.Hour)
}
sinceUnixtime := since.Unix()
untilUnixtime := until.Unix()
if sinceUnixtime >= untilUnixtime {
err = fix.Error(WrongDateRange)
return
}
var (
cumulativeTasks []*CumulativeTask
instantTasks []*InstantTask
wg = new(sync.WaitGroup)
)
// Важно вирутальная метрика зависит от реальных, которые к объекту могут быть
// привязаны или нет. Если нет - не возвращать.
resultMap = make(map[int64][]Measure)
if len(req.CumulativeMetrics) > 0 {
for _, metricID := range req.CumulativeMetrics {
task := &CumulativeTask{
In: MetricMeasuresFilter{
MetricID: metricID,
Since: sinceUnixtime,
Until: untilUnixtime,
},
WaitGroup: wg,
}
wg.Add(1)
s.cumulativeCh <- task
cumulativeTasks = append(cumulativeTasks, task)
}
}
if len(req.InstantMetrics) > 0 {
for _, metricID := range req.InstantMetrics {
task := &InstantTask{
In: MetricMeasuresFilter{
MetricID: metricID,
Since: sinceUnixtime,
Until: untilUnixtime,
},
WaitGroup: wg,
}
wg.Add(1)
s.instantCh <- task
instantTasks = append(instantTasks, task)
}
}
wg.Wait()
for _, task := range cumulativeTasks {
if task.Err != nil {
err = fmt.Errorf("get data for the metric %d: %s", task.In.MetricID, task.Err)
return
}
resultMap[task.In.MetricID] = task.Result
}
for _, task := range instantTasks {
if task.Err != nil {
err = fmt.Errorf("get data for the metric %d: %s", task.In.MetricID, task.Err)
return
}
resultMap[task.In.MetricID] = task.Result
}
return
}