package mqe // НОВАЯ ВЕРСИЯ // type MetricsDataFilter2 struct { // Since time.Time `json:"since"` // Until time.Time `json:"until"` // GroupBy string `json:"groupBy"` // InstantMetrics []InstantMetricFilter // virtualMetricID + aggregateFuncs // CumulativeMetrics []CumulativeMetricFilter // } // type InstantMetricFilter struct { // MetricID int64 // AggregateFuncs byte // LastDayOfMonth int // например, конец месяца 25 число // FirstHourOfDay int // например день начинается в 6:00 // } // type CumulativeMetricFilter struct { // MetricID int64 // LastDayOfMonth int // например, конец месяца 25 число // FirstHourOfDay int // например день начинается в 6:00 // } // func (s *MeasureQueryEngine) SelectMetricsDataInParallel2(req MetricsDataFilter2) (_ MetricsDataResultMaps, err error) { // pretty.PPrintln("SelectMetricsDataInParallel2", req) // // Строгая валидация // if req.Since.IsZero() { // err = fmt.Errorf("Since option is required") // return // } // if req.Until.IsZero() { // err = fmt.Errorf("Until option is required") // return // } // if req.Since.After(req.Until) { // err = fmt.Errorf("Wrong time range: Since %d after Until %d", req.Since.Unix(), req.Until.Unix()) // return // } // switch req.GroupBy { // case "h", "d", "m", "": // // pass // default: // err = fmt.Errorf("Wrong groupBy option value: %s", req.GroupBy) // return // } // for _, mf := range req.InstantMetrics { // if mf.MetricID == 0 { // err = fmt.Errorf("Empty instant metric's MetricID: %d", mf.MetricID) // return // } // if mf.FirstHourOfDay < 0 || mf.FirstHourOfDay > 23 { // err = fmt.Errorf("Wrong instant metric's FirstHourOfDay option value: %d", mf.FirstHourOfDay) // return // } // if mf.LastDayOfMonth < 0 || mf.LastDayOfMonth > 28 { // err = fmt.Errorf("Wrong instant metric's LastDayOfMonth option value: %d", mf.LastDayOfMonth) // return // } // } // for _, mf := range req.CumulativeMetrics { // if mf.MetricID == 0 { // err = fmt.Errorf("Empty cumulative metric's MetricID: %d", mf.MetricID) // return // } // if mf.FirstHourOfDay < 0 || mf.FirstHourOfDay > 23 { // err = fmt.Errorf("Wrong cumulative metric's FirstHourOfDay option value: %d", mf.FirstHourOfDay) // return // } // if mf.LastDayOfMonth < 0 || mf.LastDayOfMonth > 28 { // err = fmt.Errorf("Wrong cumulative metric's LastDayOfMonth option value: %d", mf.LastDayOfMonth) // return // } // } // // Корректируем, ВСЕГДА на первую и последнюю секунду суток // since := timeutil.FirstSecondInPeriod(req.Since, "d") // until := timeutil.LastSecondInPeriod(req.Until, "d") // fmt.Printf(`SelectMetricsDataInParallel2: { // Since: %s // Until: %s // GroupBy: %s // } // `, since, until, req.GroupBy) // // Важно вирутальная метрика зависит от реальных, которые к объекту могут быть // // привязаны или нет. Если нет - не возвращать. // var ( // unexpectedErrors []error // errorMutex sync.Mutex // wg sync.WaitGroup // resultMaps = MetricsDataResultMaps{ // AggregatedCumulative: make(map[int64][]CalculatedAggregatedMeasure), // AggregatedByFuncInstant: make(map[int64][]AggregatedByFuncInstantMeasure), // Cumulative: make(map[int64][]CalculatedMeasure), // Instant: make(map[int64][]CalculatedMeasure), // } // ) // wg.Add(len(req.CumulativeMetrics) + len(req.InstantMetrics)) // if req.GroupBy != "" { // // AGGREGATED // // CUMULATIVE // for _, metricFilter := range req.CumulativeMetrics { // go func(mf CumulativeMetricFilter) { // // У каждой метрики индивидуальные настройки, // // ибо метрики принадлежат разным объектам!!! // metricSince := since // metricUntil := until // if mf.FirstHourOfDay > 0 { // metricSince = since.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // metricUntil = until.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // } // result, err := s.listAggregatedCumulativeMeasures(AggregatedMeasuresFilter{ // MetricID: mf.MetricID, // Since: metricSince.Unix(), // Until: metricUntil.Unix(), // GroupBy: req.GroupBy, // LastDayOfMonth: mf.LastDayOfMonth, // FirstHourOfDay: mf.FirstHourOfDay, // }) // if err != nil { // err = fmt.Errorf("ListAggregatedCumulativeMeasures: %s", err) // errorMutex.Lock() // unexpectedErrors = append(unexpectedErrors, err) // errorMutex.Unlock() // } else { // errorMutex.Lock() // resultMaps.AggregatedCumulative[mf.MetricID] = result // errorMutex.Unlock() // } // wg.Done() // }(metricFilter) // } // // INSTANT // for _, metricFilter := range req.InstantMetrics { // go func(mf InstantMetricFilter) { // // У каждой метрики индивидуальные настройки, // // ибо метрики принадлежат разным объектам!!! // metricSince := since // metricUntil := until // if mf.FirstHourOfDay > 0 { // metricSince = since.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // metricUntil = until.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // } // result, err := s.listAggregatedByFuncsInstantMeasures(AggregatedByFuncInstantMeasuresFilter{ // MetricID: mf.MetricID, // Since: metricSince.Unix(), // Until: metricUntil.Unix(), // GroupBy: req.GroupBy, // LastDayOfMonth: mf.LastDayOfMonth, // FirstHourOfDay: mf.FirstHourOfDay, // Flags: mf.AggregateFuncs, // ФЛАГИ !!! // }) // if err != nil { // err = fmt.Errorf("ListAggregatedByFuncInstantMeasures: %s", err) // errorMutex.Lock() // unexpectedErrors = append(unexpectedErrors, err) // errorMutex.Unlock() // } else { // errorMutex.Lock() // resultMaps.AggregatedByFuncInstant[mf.MetricID] = result // errorMutex.Unlock() // } // wg.Done() // }(metricFilter) // } // } else { // // NON AGGREGATED // // CUMULATIVE // for _, metricFilter := range req.CumulativeMetrics { // go func(mf CumulativeMetricFilter) { // // У каждой метрики индивидуальные настройки, // // ибо метрики принадлежат разным объектам!!! // metricSince := since // metricUntil := until // if mf.FirstHourOfDay > 0 { // metricSince = since.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // metricUntil = until.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // } // result, err := s.listCumulativeMeasures(MeasuresFilter{ // MetricID: mf.MetricID, // Since: metricSince.Unix(), // Until: metricUntil.Unix(), // }) // if err != nil { // err = fmt.Errorf("ListCumulativeMeasures: %s", err) // errorMutex.Lock() // unexpectedErrors = append(unexpectedErrors, err) // errorMutex.Unlock() // } else { // errorMutex.Lock() // resultMaps.Cumulative[mf.MetricID] = result // errorMutex.Unlock() // } // wg.Done() // }(metricFilter) // } // // INSTANT // for _, metricFilter := range req.InstantMetrics { // go func(mf InstantMetricFilter) { // // У каждой метрики индивидуальные настройки, // // ибо метрики принадлежат разным объектам!!! // metricSince := since // metricUntil := until // if mf.FirstHourOfDay > 0 { // metricSince = since.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // metricUntil = until.Add(time.Duration(mf.FirstHourOfDay) * time.Hour) // } // result, err := s.listInstantMeasures(MeasuresFilter{ // MetricID: mf.MetricID, // Since: metricSince.Unix(), // Until: metricUntil.Unix(), // }) // if err != nil { // err = fmt.Errorf("ListInstantMeasures: %s", err) // errorMutex.Lock() // unexpectedErrors = append(unexpectedErrors, err) // errorMutex.Unlock() // } else { // errorMutex.Lock() // resultMaps.Instant[mf.MetricID] = result // errorMutex.Unlock() // } // wg.Done() // }(metricFilter) // } // } // 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 resultMaps, nil // }