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

571 lines
18 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"
"time"
)
func SplitByPeriods(producer periodProducer, calc PeriodCalculator) (result []AggregatedMeasure, err error) {
var (
isOpened bool
openedPeriod time.Time
since int64
sinceValue float64
until int64
untilValue float64
//
period time.Time
first, last PeriodBound
found bool
)
for {
period, first, last, found, err = producer.Next()
if err != nil {
return
}
if !found {
break
}
if isOpened {
if calc.IsExtendedSincePeriod(period) {
err = fmt.Errorf("extended period while already opened")
return
} else if calc.IsExtendedUntilPeriod(period) {
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: first.Time,
Values: []float64{
first.Value,
first.Value - sinceValue,
},
})
// Последний период успешно закрыли данными из Extended Until, поэтому выходим
return
} else {
if period.Equal(openedPeriod) {
// +
// Данная ситуация будет когда предыдущий период был корректно
// закрыт за 3 минуты до конца. После этого сразу открылся
// новый период. А запись для нового периода получена только
// сейчас.
if calc.IsPeriodCorrectEnds(period, last.Time) {
// CLOSE
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: last.Time,
Values: []float64{
last.Value,
last.Value - sinceValue,
},
})
// OPEN PERIOD (extended since, last 3 min closed)
openedPeriod = calc.NextPeriod(period)
if calc.IsExtendedUntilPeriod(openedPeriod) {
return
}
since = last.Time
sinceValue = last.Value
until = 0
untilValue = 0
} else {
// Период закроем где-то в будущем, на всякий случай запомнили последнее
// показание внутри периода
until = last.Time
untilValue = last.Value
}
} else {
// +
// Периоды не совпадают - смело закрываем
// Возможные ситуации:
// - был открыт мартовский период, а текущий период - апрель
// -
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: first.Time,
Values: []float64{
first.Value,
first.Value - sinceValue,
},
})
if first.Time != last.Time {
// +
openedPeriod = period
since = first.Time
sinceValue = first.Value
until = 0
untilValue = 0
if calc.IsPeriodCorrectEnds(period, last.Time) {
// Закрываем только что открытый период
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: last.Time,
Values: []float64{
last.Value,
last.Value - sinceValue,
},
})
openedPeriod = calc.NextPeriod(period)
if calc.IsExtendedUntilPeriod(openedPeriod) {
return
}
since = last.Time
sinceValue = last.Value
until = 0
untilValue = 0
} else {
// Период закроем где-то в будущем, на всякий случай запомнили последнее
// показание внутри периода
until = last.Time
untilValue = last.Value
}
} else {
// +
// first == last
if calc.IsPeriodCorrectEnds(period, last.Time) {
// Допустим закрыли предыдущий период 30 апреля в 23:58, новый стартуем за май.
openedPeriod = calc.NextPeriod(period)
if calc.IsExtendedUntilPeriod(openedPeriod) {
return
}
} else {
// Допустим закрыли предыдущий период 15 апреля, новый стартуем за апрель.
openedPeriod = period
}
since = last.Time
sinceValue = last.Value
until = 0
untilValue = 0
}
}
}
} else {
if calc.IsExtendedSincePeriod(period) {
// Если последние 3 минуты
if calc.IsPeriodCorrectEnds(period, last.Time) {
// OPEN PERIOD (extended since, last 3 min closed)
isOpened = true
openedPeriod = calc.NextPeriod(period)
since = last.Time
sinceValue = last.Value
until = 0
untilValue = 0
}
} else if calc.IsExtendedUntilPeriod(period) {
// Ничего не нашли, пустой результат
return
} else {
// OPEN PERIOD (regular first)
isOpened = true
openedPeriod = period
since = first.Time
sinceValue = first.Value
until = last.Time
untilValue = last.Value
if calc.IsPeriodCorrectEnds(openedPeriod, last.Time) {
// Закрываем только что открытый период
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: last.Time,
Values: []float64{
last.Value,
last.Value - sinceValue,
},
})
openedPeriod = calc.NextPeriod(period)
if calc.IsExtendedUntilPeriod(openedPeriod) {
return
}
since = last.Time
sinceValue = last.Value
until = 0
untilValue = 0
}
}
}
}
// Доп запрос.
if openedPeriod.IsZero() {
// Открытого периода не будет если вообще не нашли записей
return
}
next, found, err := producer.Last()
if err != nil {
return
}
if found {
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: next.Time,
Values: []float64{
next.Value,
next.Value - sinceValue,
},
})
} else {
if until > 0 {
result = append(result, AggregatedMeasure{
Period: openedPeriod.Unix(),
Since: since,
Until: until,
Values: []float64{
untilValue,
untilValue - sinceValue,
},
})
}
}
return
}
// Period - это начало периода, например для дней - это 0ч 0м 0с
// lastTm предыдущего периода всегда будет меньше current period
// since - это всегда сдвинутая метка времени (+ FirstHourOfDay * 3600)
// При FirstHourOfDay = 3 и группировке по дням, будет такая ситуация.
// Period будет указывать на 1 марта 00:00:00, а since на 1 марта 02:00:00,
// а lastTm предыдущего периода на 1 марта 01:59:00.
// Поэтому нельзя сравнивать текущий period и lastTm предыдущего периода.
// Преобразовать lastTm к period непонятно как если нужно учитывать LastDayOfMonth.
// Пример: 28 марта 00ч стало 1 апреля 00ч, а 27 марта 23:59 так и осталось.
// Можно вычислить endPeriodTime для previousPeriodEndTime, но текущий период может
// быть не строго следующим после предыдущего (например, июнь после марта).
// Правильное решение - берем firstTm текущего периода и вычисляем его реальное
// календарное начало. Например, FirstHourOfDay = 3, LastDayOfMonth = 28,
// firstTm = 1 мая. Календарное начало - 28 апреля 03:00:00. И уже с этим значением
// сравниваем previousPeriodEndTime.
// Но можно сравнивать since c lastTm и until c firstTm для определения
// extended периодов.
/*
func periodToCalendarTime(period time.Time, groupBy string, lastDayOfMonth, firstHourOfDay int) time.Time {
switch groupBy {
case "d":
if firstHourOfDay == 0 {
return period
}
// Было 0ч, стало - 2ч
return period.Add(time.Duration(firstHourOfDay) * time.Hour)
case "m":
if lastDayOfMonth == 0 {
if firstHourOfDay == 0 {
return period
}
// Было 0ч, стало - 2ч
return period.Add(time.Duration(firstHourOfDay) * time.Hour)
} else {
// текущий период начинается в прошлом месяце
tm := period.AddDate(0, -1, 0)
return time.Date(
tm.Year(),
tm.Month(),
lastDayOfMonth+1,
firstHourOfDay,
0,
0,
0,
tm.Location(),
)
}
default:
return period
}
}
*/
//var threeMinutes int64 = 3 * 60
/*
type readOptions struct {
MetricID int64 // только для readAggregatedMeasuresCanEndsInFuture
LastDayOfMonth int
FirstHourOfDay int
GroupBy string
Since int64
Until int64
ParsePeriodLayout string
}
*/
// Читает из sql.Rows показания и строит периоды
/*
func (s *MeasureQueryEngine) readAggregatedMeasuresCanEndsInFuture(tx *sql.Tx, rows *sql.Rows, opt readOptions) (result []CalculatedAggregatedMeasure, err error) {
var (
openPeriod time.Time
since int64
sinceValue float64
until int64
untilValue float64
//
frameStr string
firstTm int64
firstValue float64
lastTm int64
lastValue float64
period time.Time
)
for rows.Next() {
err = rows.Scan(&frameStr, &firstTm, &firstValue, &lastTm, &lastValue)
if err != nil {
return
}
period, err = time.ParseInLocation(opt.ParsePeriodLayout, frameStr, s.location)
if err != nil {
err = fmt.Errorf("time.ParseInLocation: %s; layout=%q; str=%q",
err, opt.ParsePeriodLayout, frameStr)
return
}
//fmt.Printf("%s s: %s, %v; u: %s, %v\n", frameStr, time.Unix(firstTm, 0), firstValue, time.Unix(lastTm, 0), lastValue)
// Читаем записи из БД. Одна запись - один период. Возможны 3 вида периодов:
// 1. extended since - период длиной до 3 минут перед since. Оптимизация, чтобы
// не отправлять дополнительный запрос.
// 2. обычный период.
// 3. extended until - период длиной до 1 часа после until. Оптимизация, чтобы
// не отправлять дополнительный запрос (в большинстве случаев).
//
// ВАЖНО!
// Любой период (1, 2 или 3) может встретится первым. Определить что это
// первая запись модно простой проверкой openPeriod.IsZero()
if lastTm < opt.Since {
// 1. extended since period
if isPeriodCorrectEnds(period, opt.GroupBy, opt.LastDayOfMonth, opt.FirstHourOfDay, lastTm) {
openPeriod = nextPeriod(period, opt.GroupBy)
since = lastTm
sinceValue = lastValue
}
continue
} else if firstTm > opt.Until {
// 3. extended until period
if openPeriod.IsZero() {
// Нет показаний кроме extended until, поэтому выходим
return
}
// закрываем открытый период
result = append(result, CalculatedAggregatedMeasure{
Period: openPeriod.Unix(),
Since: since,
Until: firstTm,
Value: firstValue,
Total: firstValue - sinceValue,
})
return
}
// 2. обычный период
if openPeriod.IsZero() {
// это первая запись (extended since period не найден)
openPeriod = period
since = firstTm
sinceValue = firstValue
} else {
if !period.Equal(openPeriod) {
//fmt.Printf("\nperiod %s != openPeriod %s\n", period, openPeriod)
// Пример: открытый период - март. В period - апрель.
// закрываем открытый период
result = append(result, CalculatedAggregatedMeasure{
Period: openPeriod.Unix(),
Since: since,
Until: firstTm,
Value: firstValue,
Total: firstValue - sinceValue,
})
// И сразу открываем новый период
openPeriod = period
since = firstTm
sinceValue = firstValue
}
// Периоды совпадают. Когда такая ситуация возможна?
// Например, на предыдущей итерации был корректно закрыт (57-59 минуты)
// прошлый период и сразу открыт новый (текущий).
}
if isPeriodCorrectEnds(openPeriod, opt.GroupBy, opt.LastDayOfMonth, opt.FirstHourOfDay, lastTm) {
// fmt.Printf("\nCORRECT END period %s\n", openPeriod)
// закрываем открытый период
result = append(result, CalculatedAggregatedMeasure{
Period: openPeriod.Unix(),
Since: since,
Until: lastTm,
Value: lastValue,
Total: lastValue - sinceValue,
})
// Открываем новый период
openPeriod = nextPeriod(openPeriod, opt.GroupBy)
since = lastTm
sinceValue = lastValue
} else {
// На последних минутах периода (57-59) не было показания, поэтому
// запоминаем конец периода какой есть. Если конец периода в будущем
// не будет найден - закроем период по текущему последнему показанию.
until = lastTm
untilValue = lastValue
}
// Валидируем новый период.
// Зачем?
// Пример: юзер задал opt.Until - 31 марта 23:59:59. Только что корректно
// закрылся март показанием от 31 марта 23:57:00. И сразу открылся новый
// период за апрель.
// Понять что закрыли последний период из запрошенных юзером просто -
// вычисляем календарное начало нового периода и сравниваем с opt.Until.
mustLessThanUntil := periodToCalendarTime(openPeriod, opt.GroupBy,
opt.LastDayOfMonth, opt.FirstHourOfDay)
if mustLessThanUntil.Unix() > opt.Until {
// Новый период вышел за opt.Until, а это означает что мы корректно
// закрыли последний период из запрошенных юзером. Bыходим.
return
}
}
// Цикл завершен, а из метода мы не вышли, а это означает одно - есть незакрытый
// период. Extended until не найден. Значит нужно отправить дополнительный запрос
// к БД.
var (
isNextValueFound bool
nextTm int64
nextValue float64
)
// Находим следующее показание метрики.
err = tx.QueryRow(`
SELECT tm, value, true
FROM f64
WHERE metricID=? AND tm > ?
ORDER BY tm ASC
LIMIT 1`,
opt.MetricID, opt.Until).Scan(&nextTm, &nextValue, &isNextValueFound)
if err != nil {
if err != sql.ErrNoRows {
return
}
// Нет записей
err = nil
}
if isNextValueFound {
// Найдено какое-то показание в будущем.
result = append(result, CalculatedAggregatedMeasure{
Period: openPeriod.Unix(),
Since: since,
Until: nextTm,
Value: nextValue,
Total: nextValue - sinceValue,
})
} else {
// Закрываем период самым последним показанием внутри периода
// (НЕ 57-59 минуты).
result = append(result, CalculatedAggregatedMeasure{
Period: openPeriod.Unix(),
Since: since,
Until: until,
Value: untilValue,
Total: untilValue - sinceValue,
})
}
return
}
*/
/*
// Читает из sql.Rows показания и строит периоды
func (s *MeasureQueryEngine) readAggregatedMeasuresEndsInsidePeriod(tx *sql.Tx, rows *sql.Rows, opt readOptions) (result []CalculatedAggregatedMeasure, err error) {
var (
previousPeriodEndTime int64
previousPeriodEndValue float64
//
frameStr string
firstTm int64
firstValue float64
lastTm int64
lastValue float64
period time.Time
)
for rows.Next() {
err = rows.Scan(&frameStr, &firstTm, &firstValue, &lastTm, &lastValue)
if err != nil {
return
}
period, err = time.ParseInLocation(opt.ParsePeriodLayout, frameStr, s.location)
if err != nil {
err = fmt.Errorf("time.ParseInLocation: %s; layout=%q; str=%q",
err, opt.ParsePeriodLayout, frameStr)
return
}
//fmt.Printf("%s s: %s, %v; u: %s, %v\n", frameStr, time.Unix(firstTm, 0), firstValue, time.Unix(lastTm, 0), lastValue)
if lastTm < opt.Since {
// extended since period
previousPeriodEndTime = lastTm
previousPeriodEndValue = lastValue
continue
} else if firstTm > opt.Until {
// extended until period
return
}
// целевой период. В будущем конец периода не ищем, поэтому закрываем сразу.
// Вычисляем реальное календарное начало периода. Например, period - 1 мая 00:00:00,
// FirstHourOfDay = 3, LastDayOfMonth = 28. Календарное начало - 29 апреля 03:00:00.
periodCalendarTime := periodToCalendarTime(period, opt.GroupBy, opt.LastDayOfMonth, opt.FirstHourOfDay)
minAge := periodCalendarTime.Add(-3 * time.Minute).Unix()
if previousPeriodEndTime > 0 && previousPeriodEndTime >= minAge {
result = append(result, CalculatedAggregatedMeasure{
Period: period.Unix(),
Since: previousPeriodEndTime,
Until: lastTm,
Value: lastValue,
Total: lastValue - previousPeriodEndValue,
})
} else {
result = append(result, CalculatedAggregatedMeasure{
Period: period.Unix(),
Since: firstTm,
Until: lastTm,
Value: lastValue,
Total: lastValue - firstValue,
})
}
previousPeriodEndTime = lastTm
previousPeriodEndValue = lastValue
}
return
}
*/