142 lines
4.0 KiB
Go
142 lines
4.0 KiB
Go
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
|
||
}
|