package mqe import ( "database/sql" "fmt" "sync" "time" "gordenko.dev/dima/pretty" "gordenko.dev/dima/qb/timeutil" ) //////////////////////////// type MetricMeasure struct { Time int64 `json:"t"` Value float64 `json:"v"` } type MetricError struct { Time int64 `json:"time"` Code int64 `json:"code"` } type ReadingsFilter struct { Since string `json:"since"` Until string `json:"until"` MetricIDs []int64 } type MetricReadings struct { MetricID int64 Measures []MetricMeasure Errors []MetricError } // SelectReadingsInParallel - метод достает из БД данные всех метрик, которые связаны // с объектом. Учитываются настройки объекта (LastDayOfMonth, FirstHourOfDay), а также // фильтры выбранные пользователем на странице объекта (Since, Until, GroupBy). // ВАЖНО! // Запросы по каждой метрике отправляются в БД в отдельной горутине, что позволяет // в разы ускорить получение данных по объекту. func (s *MeasureQueryEngine) SelectReadingsInParallel(req ReadingsFilter) (result []MetricReadings, err error) { pretty.PPrintln("SelectReadingsInParallel", req) if len(req.MetricIDs) == 0 { return } 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)) fmt.Printf(`SelectReadingsInParallel: { Since: %s Until: %s MetricIDs: %v } `, since, until, req.MetricIDs) // FIX проверить как добавляются часы в DST часовых поясах // Для req.FirstHourOfDay=4 since должен быть 4:00:00, until должен быть // 3:59:59 следующего дня sinceUnixtime := since.Unix() untilUnixtime := until.Unix() if sinceUnixtime >= untilUnixtime { err = fix.Error(WrongDateRange) return } // Важно вирутальная метрика зависит от реальных, которые к объекту могут быть // привязаны или нет. Если нет - не возвращать. var ( tasks []*MetricReadingsTask wg = new(sync.WaitGroup) ) for _, metricID := range req.MetricIDs { task := &MetricReadingsTask{ In: MetricMeasuresFilter{ MetricID: metricID, Since: sinceUnixtime, Until: untilUnixtime, }, WaitGroup: wg, } wg.Add(1) s.metricReadingsCh <- 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 } result = append(result, task.Result) } return } // ListInstantMeasures - cписок показаний для юнитов без накопительного итога // (например, Давление) func (s *MeasureQueryEngine) listMetricReadings(req MetricMeasuresFilter) (result MetricReadings, err error) { result.MetricID = req.MetricID // buf, _ := json.MarshalIndent(req, "", " ") // fmt.Printf("listInstantMeasures: %s\n\n", buf) tx, err := s.db.Driver().Begin() if err != nil { return } defer tx.Rollback() rows, err := tx.Query(` SELECT tm, value FROM f64 WHERE metricID=? AND tm BETWEEN ? AND ? ORDER BY tm ASC`, req.MetricID, req.Since, req.Until) if err != nil { if err == sql.ErrNoRows { err = nil } return } defer rows.Close() for rows.Next() { var measure MetricMeasure err = rows.Scan(&measure.Time, &measure.Value) if err != nil { return } // fix - для cumulative метрик учет поправок result.Measures = append(result.Measures, measure) } if err = rows.Err(); err != nil { return } //for _, m := range result { // fmt.Printf("raw: %d\t%f\n", m.Time, m.Value) //} //fmt.Printf("listInstantMeasures (qty): %d\n", len(result)) return }