139 lines
4.3 KiB
Go
139 lines
4.3 KiB
Go
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
|
||
}
|