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 }