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
|
|||
|
|
}
|