package mqe import ( "fmt" "strings" "sync" "time" "gordenko.dev/dima/qb/timeutil" ) // TOTALS type TotalMetricSpec struct { MetricID int64 LastDayOfMonth int FirstHourOfDay int } // LastDayOfMonth и FirstHourOfDay нет, ибо у каждой метрики свои настройки type RangeTotalsFilter struct { Since string Until string // Если учитывать не нужно - программа уровнем выше может передать нули. // FIX у каждой метрики свои настройки // искать ли конец периода в следующих периодах, если в последние 3 минуты периода // не было показаний CumulativeMetrics []TotalMetricSpec // metricID } // SelectRangeTotalsInParallel - метод рассчитывает накопленный объем от Since до Until // по каждой метрике из массива CumulativeMetricIDs. Учитываются настройки объекта // (LastDayOfMonth, FirstHourOfDay). // ВАЖНО! // Запросы по каждой метрике отправляются в БД в отдельной горутине, что позволяет // в разы ускорить получение данных по объекту. func (s *MeasureQueryEngine) GetRangeTotalsInParallel(req RangeTotalsFilter) (resultMap map[int64]RangeTotal, err error) { if req.Since == "" { err = fix.Field(EmptyValue, "since") return } if req.Until == "" { err = fix.Field(EmptyValue, "until") return } // ВАЖНО! // Если в строке нет таймзоны - Parse функция возвращает время в UTC. // Поэтому нужно использовать ParseInLocation! since, err := time.ParseInLocation("2006-01-02", req.Since, s.location) if err != nil { err = fix.Field(WrongValue, "since") return } until, err := time.ParseInLocation("2006-01-02", req.Until, s.location) if err != nil { err = fix.Field(WrongValue, "until") return } // Корректируем, чтобы указывало на последнюю секунду суток until = timeutil.LastSecondInPeriod(until, string(ByDay)) var metrics []string for _, metric := range req.CumulativeMetrics { metrics = append(metrics, fmt.Sprintf(` { MetricID: %d LastDayOfMonth: %d FirstHourOfDay: %d }`, metric.MetricID, metric.LastDayOfMonth, metric.FirstHourOfDay)) } fmt.Printf(`GetRangeTotalsInParallel: { Since: %s Until: %s CumulativeMetrics [ %s ] } `, since, until, strings.Join(metrics, "\n")) if since.Unix() >= until.Unix() { err = fix.Error(WrongDateRange) return } var ( tasks []*CumulativeTotalTask wg = new(sync.WaitGroup) ) resultMap = make(map[int64]RangeTotal) // CUMULATIVE METRIC TOTALS if len(req.CumulativeMetrics) > 0 { for _, metric := range req.CumulativeMetrics { // У каждой метрики индивидуальные настройки, // ибо метрики принадлежат разным объектам!!! metricSince := since metricUntil := until if metric.FirstHourOfDay > 0 { metricSince = since.Add(time.Duration(metric.FirstHourOfDay) * time.Hour) metricUntil = until.Add(time.Duration(metric.FirstHourOfDay) * time.Hour) } task := &CumulativeTotalTask{ In: getCumulativeTotalFilter{ MetricID: metric.MetricID, Since: metricSince.Unix(), Until: metricUntil.Unix(), LastDayOfMonth: metric.LastDayOfMonth, FirstHourOfDay: metric.FirstHourOfDay, }, WaitGroup: wg, } wg.Add(1) s.cumulativeTotalCh <- task tasks = append(tasks, task) } } wg.Wait() for _, task := range tasks { if task.Err != nil { err = fmt.Errorf("get data for the metric %d: %s", task.In.MetricID, task.Err) return } if task.Result != nil { resultMap[task.In.MetricID] = *task.Result } else { resultMap[task.In.MetricID] = RangeTotal{} } } return }