586 lines
16 KiB
Go
586 lines
16 KiB
Go
package mqe
|
||
|
||
import (
|
||
"database/sql"
|
||
"fmt"
|
||
"time"
|
||
)
|
||
|
||
// func applyCorrections(corrections []model.F64Correction, measures []Measure) []Measure {
|
||
// for idx, measure := range measures {
|
||
// var isCorrected bool
|
||
// for _, correction := range corrections {
|
||
// if measure.Time >= correction.Time {
|
||
// measure.Value += correction.Value
|
||
// isCorrected = true
|
||
// }
|
||
// }
|
||
// if isCorrected {
|
||
// measures[idx] = measure
|
||
// }
|
||
// }
|
||
// return measures
|
||
// }
|
||
|
||
func applyCorrectionsToMeasure(corrections []_f64Correction, tm int64, value float64) float64 {
|
||
for _, correction := range corrections {
|
||
if tm >= correction.Time {
|
||
value += correction.Value
|
||
}
|
||
}
|
||
return value
|
||
}
|
||
|
||
func applyCorrectionsToMeasures(corrections []_f64Correction, measures []_measure) {
|
||
for idx, m := range measures {
|
||
for _, correction := range corrections {
|
||
if m.Time >= correction.Time {
|
||
m.Value += correction.Value
|
||
}
|
||
}
|
||
measures[idx] = m
|
||
}
|
||
}
|
||
|
||
func applyCorrectionsToRawMeasures(corrections []_f64Correction, measures []RawMeasure) {
|
||
for idx, m := range measures {
|
||
for _, correction := range corrections {
|
||
if m.Time >= correction.Time {
|
||
m.Value += correction.Value
|
||
}
|
||
}
|
||
measures[idx] = m
|
||
}
|
||
}
|
||
|
||
type _measure struct {
|
||
Time int64
|
||
Value float64
|
||
}
|
||
|
||
// ВАЖНО!
|
||
// ВО ВСЕХ ЗАПРОСАХ ОБЯЗАТЕЛЬНА КОРРЕКТНАЯ СОРТИРОВКА!
|
||
|
||
type _f64Correction struct {
|
||
Time int64
|
||
Value float64
|
||
}
|
||
|
||
func listF64CorrectionsTx(tx *sql.Tx, metricID int64) (list []_f64Correction, err error) {
|
||
rows, err := tx.Query(
|
||
"SELECT tm, value FROM metric_corrections WHERE metricID=? ORDER BY tm ASC",
|
||
metricID)
|
||
if err != nil {
|
||
if err == sql.ErrNoRows {
|
||
err = nil
|
||
}
|
||
return
|
||
}
|
||
defer rows.Close()
|
||
|
||
for rows.Next() {
|
||
var correction _f64Correction
|
||
err = rows.Scan(&correction.Time, &correction.Value)
|
||
if err != nil {
|
||
return
|
||
}
|
||
list = append(list, correction)
|
||
}
|
||
return
|
||
}
|
||
|
||
// result = append(result, Measure{
|
||
// Time: tm,
|
||
// Values: []float64{
|
||
// value,
|
||
// value - prevValue,
|
||
// },
|
||
// })
|
||
|
||
// listCumulativeMeasures - возвращает список CumulativeMeasure.
|
||
// Метод ничего не знает про FirstHourOfDay, но извне границы Since и Until могут быть
|
||
// скорректированы с учетом FirstHourOfDay. Это не поломает выборку.
|
||
func (s *MeasureQueryEngine) listCumulativeMeasures(req MetricMeasuresFilter) (_ []Measure, err error) {
|
||
tx, err := s.db.Driver().Begin()
|
||
if err != nil {
|
||
return
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
// corrections, err := listF64CorrectionsTx(tx, req.MetricID)
|
||
// if err != nil {
|
||
// return
|
||
// }
|
||
|
||
// ВАЖНО!
|
||
// Оптимизация! Запрашиваем больше данных чем нужно, чтобы посчитать Total для
|
||
// самого старого измерения в result.
|
||
// Отнимаем от since еще 1 час (с запасом) чтобы зацепить extendedMeasure
|
||
extendedSince := req.Since - 3600
|
||
|
||
rows, err := tx.Query(`
|
||
SELECT tm, value
|
||
FROM f64
|
||
WHERE metricID=? AND tm BETWEEN ? AND ?
|
||
ORDER BY tm ASC`,
|
||
req.MetricID, extendedSince, req.Until)
|
||
if err != nil {
|
||
if err == sql.ErrNoRows {
|
||
err = nil
|
||
}
|
||
return
|
||
}
|
||
defer rows.Close()
|
||
|
||
var (
|
||
result []Measure
|
||
measures []_measure
|
||
prev *_measure
|
||
)
|
||
|
||
// Вычитываем все найденные показания
|
||
for rows.Next() {
|
||
var m _measure
|
||
err = rows.Scan(&m.Time, &m.Value)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if m.Time < req.Since {
|
||
// пока не найдено первое показание внутри диапазона обновлем extendedValue
|
||
prev = &_measure{
|
||
Time: m.Time,
|
||
Value: m.Value,
|
||
}
|
||
} else {
|
||
measures = append(measures, m)
|
||
break
|
||
}
|
||
}
|
||
|
||
// Остальные показания вычитываем без проверок
|
||
for rows.Next() {
|
||
var m _measure
|
||
err = rows.Scan(&m.Time, &m.Value)
|
||
if err != nil {
|
||
return
|
||
}
|
||
measures = append(measures, m)
|
||
}
|
||
|
||
if err = rows.Err(); err != nil {
|
||
return
|
||
}
|
||
|
||
if len(measures) == 0 {
|
||
return
|
||
}
|
||
|
||
if prev == nil {
|
||
// extended not found
|
||
// extended показание не найдено, а значит для первого показания не сможем показать total.
|
||
// Нужен дополнительный запрос.
|
||
// Находим максимально свежее показание метрики до extendedSince, от которого и будем
|
||
// считать total.
|
||
var (
|
||
tm int64
|
||
value float64
|
||
)
|
||
|
||
err = tx.QueryRow(`
|
||
SELECT tm, value
|
||
FROM f64
|
||
WHERE metricID=? AND tm < ?
|
||
ORDER BY tm DESC
|
||
LIMIT 1`,
|
||
req.MetricID, extendedSince).Scan(&tm, &value)
|
||
|
||
if err != nil {
|
||
if err != sql.ErrNoRows {
|
||
return
|
||
}
|
||
err = nil
|
||
} else {
|
||
prev = &_measure{
|
||
Time: tm,
|
||
Value: value,
|
||
}
|
||
}
|
||
}
|
||
|
||
// ВАЖНО!
|
||
// Сперва применяем коррекции к показаниям, а затем рассчитываем сумму за период
|
||
|
||
// if len(corrections) > 0 {
|
||
// applyCorrectionsToMeasures(corrections, measures)
|
||
// if prev != nil {
|
||
// prev.Value = applyCorrectionsToMeasure(corrections, prev.Time, prev.Value)
|
||
// }
|
||
// }
|
||
|
||
var prevValue float64
|
||
if prev != nil {
|
||
prevValue = prev.Value
|
||
}
|
||
|
||
for _, m := range measures {
|
||
result = append(result, Measure{
|
||
Time: m.Time,
|
||
Values: []float64{
|
||
m.Value,
|
||
m.Value - prevValue,
|
||
},
|
||
})
|
||
prevValue = m.Value
|
||
}
|
||
|
||
// Корректируем самый первый тотал. Eсли prev показания не было -> total = 0
|
||
if prev == nil {
|
||
first := result[0]
|
||
first.Values[1] = 0
|
||
result[0] = first
|
||
}
|
||
|
||
return result, nil
|
||
}
|
||
|
||
type AggregatedMeasuresFilter struct {
|
||
MetricID int64 `json:"metricID"`
|
||
Since int64 `json:"since"`
|
||
Until int64 `json:"until"`
|
||
GroupBy GroupBy `json:"groupBy"`
|
||
LastDayOfMonth int `json:"lastDayOfMonth"` // например, конец месяца 25 число
|
||
FirstHourOfDay int `json:"firstHourOfDay"` // например день начинается в 6:00
|
||
}
|
||
|
||
// listAggregatedCumulativeMeasures - возвращает список AggregatedCumulativeMeasure,
|
||
// сгрупированный по какому-то периоду
|
||
func (s *MeasureQueryEngine) listAggregatedCumulativeMeasures(req AggregatedMeasuresFilter) (_ []AggregatedMeasure, err error) {
|
||
//pretty.PPrintln("listAggregatedCumulativeMeasures", req)
|
||
|
||
tx, err := s.db.Driver().Begin()
|
||
if err != nil {
|
||
return
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
corrections, err := listF64CorrectionsTx(tx, req.MetricID)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
var (
|
||
// чтобы не экранировать символы в шаблоне строки
|
||
groupByFormat string
|
||
parsePeriodLayout string
|
||
)
|
||
|
||
switch req.GroupBy {
|
||
case ByHour:
|
||
groupByFormat = "%Y%m%d%H"
|
||
parsePeriodLayout = hourPeriodLayout
|
||
|
||
case ByDay:
|
||
groupByFormat = "%Y%m%d"
|
||
parsePeriodLayout = dayPeriodLayout
|
||
|
||
case ByMonth:
|
||
groupByFormat = "%Y%m"
|
||
parsePeriodLayout = monthPeriodLayout
|
||
|
||
default:
|
||
// Защита от дурака
|
||
err = fmt.Errorf("unknown groupBy: %s", req.GroupBy)
|
||
return
|
||
}
|
||
|
||
// ВАЖНО!
|
||
// Оптимизация! Запрашиваем больше данных чем нужно, чтобы посчитать начало
|
||
// первого и конец последнего периодов без дополнительных запросов.
|
||
//
|
||
|
||
// Отнимаем 3 минуты чтобы зацепить последнее показание предыдущего периода
|
||
extendedSince := req.Since - 3*60
|
||
// Добавляем 1 час, чтобы зацепить первое показание следующего периода
|
||
extendedUntil := req.Until + 60*60
|
||
|
||
//fmt.Printf("extendedSince: %s\n\n", extendedSince)
|
||
|
||
// Для каждого периода выбираем первое и последнее (по времени) измерение.
|
||
// Если делать выборку одним запросом без UNION, а именно сделать
|
||
// INNER JOIN ... ON f64.tm=x.maxtm OR f64.tm=x.mintm - все начинает тормозить.
|
||
// Два запроса работают достаточно быстро.
|
||
|
||
// ВАЖНО!
|
||
// Не забыть про сортировку tm DESC!
|
||
|
||
dt, err := getPeriodDateTime(req.GroupBy, req.LastDayOfMonth, req.FirstHourOfDay)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
//fmt.Printf("DT: %s\n\n", dt)
|
||
|
||
query := fmt.Sprintf(`
|
||
WITH
|
||
measures AS (SELECT * FROM f64 WHERE metricID=? AND tm BETWEEN ? AND ?),
|
||
aggregated AS (
|
||
SELECT MIN(tm) as minTm, MAX(tm) as maxTm, DATE_FORMAT(%s, '%s') as periodStr
|
||
FROM measures
|
||
GROUP BY periodStr
|
||
)
|
||
SELECT aggregated.periodStr, firstMeasures.tm, firstMeasures.value, lastMeasures.tm, lastMeasures.value
|
||
FROM aggregated
|
||
INNER JOIN measures firstMeasures ON firstMeasures.tm = aggregated.minTm
|
||
INNER JOIN measures lastMeasures ON lastMeasures.tm = aggregated.maxTm
|
||
ORDER BY aggregated.periodStr ASC`,
|
||
dt, groupByFormat)
|
||
|
||
rows, err := tx.Query(query, req.MetricID, extendedSince, extendedUntil)
|
||
if err != nil {
|
||
if err == sql.ErrNoRows {
|
||
err = nil
|
||
}
|
||
return
|
||
}
|
||
defer rows.Close()
|
||
|
||
var result []AggregatedMeasure
|
||
|
||
periodCalculator := NewPeriodCalculator(PeriodCalculatorOptions{
|
||
GroupBy: req.GroupBy,
|
||
LastDayOfMonth: req.LastDayOfMonth,
|
||
FirstHourOfDay: req.FirstHourOfDay,
|
||
Since: time.Unix(req.Since, 0),
|
||
Until: time.Unix(req.Until, 0),
|
||
})
|
||
|
||
periodProducer := NewPeriodProducer(PeriodProducerOptions{
|
||
MetricID: req.MetricID,
|
||
ParsePeriodLayout: parsePeriodLayout,
|
||
Location: s.location,
|
||
Until: req.Until,
|
||
Corrections: corrections,
|
||
Rows: rows,
|
||
Tx: tx,
|
||
})
|
||
|
||
result, err = SplitByPeriods(periodProducer, periodCalculator)
|
||
if err != nil {
|
||
return
|
||
}
|
||
return result, nil
|
||
}
|
||
|
||
type RangeTotal struct {
|
||
Since int64 `json:"since"` // реальная граница
|
||
Until int64 `json:"until"` // реальная граница
|
||
StartValue float64 `json:"startValue"` // на конец периода
|
||
EndValue float64 `json:"endValue"` // на конец периода
|
||
}
|
||
|
||
type getCumulativeTotalFilter struct {
|
||
MetricID int64
|
||
Since int64
|
||
Until int64
|
||
LastDayOfMonth int
|
||
FirstHourOfDay int
|
||
}
|
||
|
||
func (s *MeasureQueryEngine) getCumulativeTotal(req getCumulativeTotalFilter) (_ *RangeTotal, err error) {
|
||
measures, err := s.listAggregatedCumulativeMeasures(AggregatedMeasuresFilter{
|
||
MetricID: req.MetricID,
|
||
Since: req.Since,
|
||
Until: req.Until,
|
||
GroupBy: ByDay,
|
||
LastDayOfMonth: req.LastDayOfMonth,
|
||
FirstHourOfDay: req.FirstHourOfDay,
|
||
})
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if len(measures) == 0 {
|
||
return
|
||
}
|
||
|
||
//pretty.PPrintln("TOTAL", measures)
|
||
|
||
first := measures[0]
|
||
last := measures[len(measures)-1]
|
||
|
||
return &RangeTotal{
|
||
Since: first.Since,
|
||
Until: last.Until,
|
||
StartValue: first.Values[0] - first.Values[1], // value on period end - period total
|
||
EndValue: last.Values[0],
|
||
}, nil
|
||
}
|
||
|
||
/*
|
||
// GetCumulativeTotal - фикс, добавить опции FirstHourInDay, IsFindPeriodEndsInFuture
|
||
func (s *MeasureQueryEngine) getCumulativeTotal(req getCumulativeTotalFilter) (_ *RangeTotal, err error) {
|
||
tx, err := s.db.Driver().Begin()
|
||
if err != nil {
|
||
return
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
var (
|
||
firstTm, lastTm, futureTm int64
|
||
firstValue, lastValue, futureValue float64
|
||
foundInPast, foundInFuture, foundFirstInRange, foundLastInRange bool
|
||
)
|
||
|
||
// Поиск первого показания
|
||
|
||
// Заглядываем в прошлое на 3 минуты
|
||
firstTm, firstValue, foundInPast, err = s.findLastInRange(tx, req.MetricID, req.Since-3*60, req.Since)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if !foundInPast {
|
||
//
|
||
firstTm, firstValue, foundFirstInRange, err = s.findFirstInRange(tx, req.MetricID, req.Since, req.Until)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if !foundFirstInRange {
|
||
// Начало периода не найдено - значит нет показаний для указанного периода
|
||
return
|
||
}
|
||
}
|
||
|
||
// Поиск последнего показания
|
||
|
||
lastTm, lastValue, foundLastInRange, err = s.findLastInRange(tx, req.MetricID, req.Since, req.Until)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if req.CanPeriodEndsInFuture {
|
||
// Если можем искать конец в периода в будущем
|
||
if foundLastInRange {
|
||
if lastTm > (req.Until - 3*60) {
|
||
// Если lastTm в последние 3 минуты периода - период корректно завершен
|
||
return &RangeTotal{
|
||
Since: firstTm,
|
||
Until: lastTm,
|
||
Value: lastValue,
|
||
Total: lastValue - firstValue,
|
||
}, nil
|
||
}
|
||
}
|
||
// Период не завершен в последние 3 минуты периода, поэтому ищем в будущем
|
||
futureTm, futureValue, foundInFuture, err = s.findFirstAfter(tx, req.MetricID, req.Until)
|
||
if err != nil {
|
||
return
|
||
}
|
||
|
||
if foundInFuture {
|
||
return &RangeTotal{
|
||
Since: firstTm,
|
||
Until: futureTm,
|
||
Value: futureValue,
|
||
Total: futureValue - firstValue,
|
||
}, nil
|
||
} else {
|
||
// Супер важная проверка, ибо возможна ситуация, когда когда нашли начало
|
||
// текущего периода в конце предыдущего (foundInPast). А это показание было
|
||
// последним в БД.
|
||
if foundLastInRange {
|
||
// Последнее показание внутри периода нашли, неважно когда оно пришло - нам
|
||
// подходит.
|
||
return &RangeTotal{
|
||
Since: firstTm,
|
||
Until: lastTm,
|
||
Value: lastValue,
|
||
Total: lastValue - firstValue,
|
||
}, nil
|
||
}
|
||
// Последнее показание внутри периода не найдено, а значит их не было вообще.
|
||
// Просто выходим
|
||
return
|
||
}
|
||
} else {
|
||
// В будущем искать нельзя
|
||
|
||
// Супер важная проверка, ибо возможна ситуация, когда когда нашли начало
|
||
// текущего периода в конце предыдущего (foundInPast). А это показание было
|
||
// последним в БД.
|
||
if foundLastInRange {
|
||
// Последнее показание внутри периода нашли, неважно когда оно пришло - нам
|
||
// подходит.
|
||
return &RangeTotal{
|
||
Since: firstTm,
|
||
Until: lastTm,
|
||
Value: lastValue,
|
||
Total: lastValue - firstValue,
|
||
}, nil
|
||
}
|
||
// Последнее показание внутри периода не найдено, а значит их не было вообще.
|
||
// Просто выходим
|
||
return
|
||
}
|
||
}
|
||
*/
|
||
|
||
func (s *MeasureQueryEngine) findLastInRangeTx(tx *sql.Tx, metricID int64, since, until int64) (tm int64, value float64, found bool, err error) {
|
||
err = tx.QueryRow(`
|
||
SELECT tm, value, true
|
||
FROM f64
|
||
WHERE metricID=? AND tm >= ? AND tm < ?
|
||
ORDER BY tm DESC
|
||
LIMIT 1`,
|
||
metricID, since, until).Scan(&tm, &value, &found)
|
||
|
||
if err != nil {
|
||
if err != sql.ErrNoRows {
|
||
return
|
||
}
|
||
// Нет записей
|
||
err = nil
|
||
}
|
||
return
|
||
}
|
||
|
||
func (s *MeasureQueryEngine) findFirstInRangeTx(tx *sql.Tx, metricID int64, since, until int64) (tm int64, value float64, found bool, err error) {
|
||
err = tx.QueryRow(`
|
||
SELECT tm, value, true
|
||
FROM f64
|
||
WHERE metricID=? AND tm >= ? AND tm < ?
|
||
ORDER BY tm ASC
|
||
LIMIT 1`,
|
||
metricID, since, until).Scan(&tm, &value, &found)
|
||
|
||
if err != nil {
|
||
if err != sql.ErrNoRows {
|
||
return
|
||
}
|
||
// Нет записей
|
||
err = nil
|
||
}
|
||
return
|
||
}
|
||
|
||
func (s *MeasureQueryEngine) findFirstAfterTx(tx *sql.Tx, metricID int64, since int64) (tm int64, value float64, found bool, err error) {
|
||
err = tx.QueryRow(`
|
||
SELECT tm, value, true
|
||
FROM f64
|
||
WHERE metricID=? AND tm >= ?
|
||
ORDER BY tm ASC
|
||
LIMIT 1`,
|
||
metricID, since).Scan(&tm, &value, &found)
|
||
|
||
if err != nil {
|
||
if err != sql.ErrNoRows {
|
||
return
|
||
}
|
||
// Нет записей
|
||
err = nil
|
||
}
|
||
return
|
||
}
|