Files
qb/mqe/aggregated.go
2026-06-19 07:25:46 +03:00

341 lines
9.8 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package mqe
import (
"fmt"
"strings"
"sync"
"time"
"gordenko.dev/dima/qb/timeutil"
)
type AggregatedMeasure struct {
Period int64 `json:"p"`
Since int64 `json:"s"`
Until int64 `json:"u"`
// для cumulative метрики 2 значения - value и total
// для instant метрики от 1 до 3 значений - min, max, avg
Values []float64 `json:"v"`
}
type AggregatedFilter struct {
Since time.Time
Until time.Time
GroupBy GroupBy
LastDayOfMonth int // например, конец месяца 25 число
FirstHourOfDay int // например день начинается в 6:00
InstantMetrics []InstantMetricAndFuncs
CumulativeMetrics []int64
}
type InstantMetricAndFuncs struct {
MetricID int64
AggregateFuncs byte
}
func (s *MeasureQueryEngine) SelectAggregatedData(req AggregatedFilter) (resultMap map[int64][]AggregatedMeasure, err error) {
since := req.Since
until := req.Until
// Валидация
switch req.GroupBy {
case ByHour, ByDay, ByMonth:
// pass
default:
err = fix.Field(WrongValue, "groupBy")
return
}
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
}
// Корректируем, ВСЕГДА на последнюю секунду суток
//until = timeutil.LastSecondInPeriod(until, "d")
// fmt.Printf(`SelectAggregatedData: {
// Since: %s
// Until: %s
// GroupBy: %s
// FirstHourOfDay: %d
// LastDayOfMonth: %d
// }
// `, since, until, req.GroupBy, 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 []*AggregatedCumulativeTask
instantTasks []*AggregatedInstantTask
wg = new(sync.WaitGroup)
)
// Важно вирутальная метрика зависит от реальных, которые к объекту могут быть
// привязаны или нет. Если нет - не возвращать.
resultMap = make(map[int64][]AggregatedMeasure)
if len(req.CumulativeMetrics) > 0 {
for _, metricID := range req.CumulativeMetrics {
task := &AggregatedCumulativeTask{
In: AggregatedMeasuresFilter{
MetricID: metricID,
Since: sinceUnixtime,
Until: untilUnixtime,
GroupBy: req.GroupBy,
LastDayOfMonth: req.LastDayOfMonth,
FirstHourOfDay: req.FirstHourOfDay,
},
WaitGroup: wg,
}
wg.Add(1)
s.aggregatedCumulativeCh <- task
cumulativeTasks = append(cumulativeTasks, task)
}
}
if len(req.InstantMetrics) > 0 {
for _, metric := range req.InstantMetrics {
task := &AggregatedInstantTask{
In: AggregatedInstantMeasuresIn{
MetricID: metric.MetricID,
Since: sinceUnixtime,
Until: untilUnixtime,
GroupBy: req.GroupBy,
LastDayOfMonth: req.LastDayOfMonth,
FirstHourOfDay: req.FirstHourOfDay,
Flags: metric.AggregateFuncs, // ФЛАГИ !!!
},
WaitGroup: wg,
}
wg.Add(1)
s.aggregatedInstantCh <- 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
}
type IndividualAggregatedFilter struct {
Since time.Time
Until time.Time
GroupBy GroupBy
InstantMetrics []IndividualInstantMetric
CumulativeMetrics []IndividualCumulativeMetric
}
type IndividualInstantMetric struct {
MetricID int64
AggregateFuncs byte
LastDayOfMonth int // например, конец месяца 25 число
FirstHourOfDay int // например день начинается в 6:00
}
type IndividualCumulativeMetric struct {
MetricID int64
LastDayOfMonth int // например, конец месяца 25 число
FirstHourOfDay int // например день начинается в 6:00
}
func (s *MeasureQueryEngine) SelectIndividualAggregatedData(req IndividualAggregatedFilter) (resultMap map[int64][]AggregatedMeasure, err error) {
if req.Since.IsZero() {
err = fmt.Errorf("zero Since")
return
}
if req.Until.IsZero() {
err = fmt.Errorf("zero Until")
return
}
if req.Since.After(req.Until) {
err = fmt.Errorf("wrong time range: since %s after until %s", req.Since, req.Until)
return
}
// Корректируем, ВСЕГДА на первую и последнюю секунду суток
req.Since = timeutil.FirstSecondInPeriod(req.Since, "d")
req.Until = timeutil.LastSecondInPeriod(req.Until, "d")
// Валидация
switch req.GroupBy {
case ByHour, ByDay, ByMonth:
// pass
default:
err = fix.Field(WrongValue, "groupBy")
return
}
// fmt.Printf(`SelectIndividualAggregatedData: {
// Since: %s
// Until: %s
// GroupBy: %s
// }
// `, req.Since, req.Until, req.GroupBy)
// Важно вирутальная метрика зависит от реальных, которые к объекту могут быть
// привязаны или нет. Если нет - не возвращать.
resultMap = make(map[int64][]AggregatedMeasure)
var (
unexpectedErrors []error
errorMutex sync.Mutex
wg sync.WaitGroup
)
wg.Add(len(req.CumulativeMetrics) + len(req.InstantMetrics))
if len(req.CumulativeMetrics) > 0 {
for _, m := range req.CumulativeMetrics {
go func(metric IndividualCumulativeMetric) {
if metric.FirstHourOfDay < 0 || metric.FirstHourOfDay > 23 {
err = fix.Field(InvalidFirstHourOfDay, "firstHourOfDay")
return
}
if metric.LastDayOfMonth < 0 || metric.LastDayOfMonth > 28 {
err = fix.Field(InvalidLastDayOfMonth, "lastDayOfMonth")
return
}
since := req.Since
until := req.Until
// FIX проверить как добавляются часы в DST часовых поясах
// Для req.FirstHourOfDay=4 since должен быть 4:00:00, until должен быть
// 3:59:59 следующего дня
if metric.FirstHourOfDay > 0 {
since = req.Since.Add(time.Duration(metric.FirstHourOfDay) * time.Hour)
until = req.Until.Add(time.Duration(metric.FirstHourOfDay) * time.Hour)
}
result, err := s.listAggregatedCumulativeMeasures(AggregatedMeasuresFilter{
MetricID: metric.MetricID,
Since: since.Unix(),
Until: until.Unix(),
GroupBy: req.GroupBy,
LastDayOfMonth: metric.LastDayOfMonth,
FirstHourOfDay: metric.FirstHourOfDay,
})
if err != nil {
err = fmt.Errorf("ListAggregatedCumulativeMeasures: %s", err)
errorMutex.Lock()
unexpectedErrors = append(unexpectedErrors, err)
errorMutex.Unlock()
} else {
errorMutex.Lock()
resultMap[metric.MetricID] = result
errorMutex.Unlock()
}
wg.Done()
}(m)
}
}
if len(req.InstantMetrics) > 0 {
for _, m := range req.InstantMetrics {
go func(metric IndividualInstantMetric) {
if metric.FirstHourOfDay < 0 || metric.FirstHourOfDay > 23 {
err = fix.Field(InvalidFirstHourOfDay, "firstHourOfDay")
return
}
if metric.LastDayOfMonth < 0 || metric.LastDayOfMonth > 28 {
err = fix.Field(InvalidLastDayOfMonth, "lastDayOfMonth")
return
}
since := req.Since
until := req.Until
// FIX проверить как добавляются часы в DST часовых поясах
// Для req.FirstHourOfDay=4 since должен быть 4:00:00, until должен быть
// 3:59:59 следующего дня
if metric.FirstHourOfDay > 0 {
since = req.Since.Add(time.Duration(metric.FirstHourOfDay) * time.Hour)
until = req.Until.Add(time.Duration(metric.FirstHourOfDay) * time.Hour)
}
result, err := s.listAggregatedInstantMeasures(AggregatedInstantMeasuresIn{
MetricID: metric.MetricID,
Since: since.Unix(),
Until: until.Unix(),
GroupBy: req.GroupBy,
LastDayOfMonth: metric.LastDayOfMonth,
FirstHourOfDay: metric.FirstHourOfDay,
Flags: metric.AggregateFuncs, // ФЛАГИ !!!
})
if err != nil {
err = fmt.Errorf("listAggregatedInstantMeasures: %s", err)
errorMutex.Lock()
unexpectedErrors = append(unexpectedErrors, err)
errorMutex.Unlock()
} else {
errorMutex.Lock()
resultMap[metric.MetricID] = result
errorMutex.Unlock()
}
wg.Done()
}(m)
}
}
wg.Wait()
// Проверяем что все запросы завершились удачно
if len(unexpectedErrors) > 0 {
// Нумеруем сообщения об ошибках и склеиваем в одно большое сообщение
var errorStrings []string
for idx, err := range unexpectedErrors {
errorStrings = append(errorStrings, fmt.Sprintf("\n#%d %s", idx+1, err))
}
err = fmt.Errorf("%d errors occured:%s", len(unexpectedErrors),
strings.Join(errorStrings, ""))
return
}
return
}