From b085e70957d99de2abe5deb5d99d404fd02b0c72 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Tue, 21 Nov 2023 17:44:22 +0800 Subject: [PATCH] fix(monitor): query multiple fileds from victoria-metrics --- .../driver/victoriametrics/result_parser.go | 29 +++++++++++++++++-- pkg/monitor/tsdb/driver/victoriametrics/vm.go | 2 +- 2 files changed, 27 insertions(+), 4 deletions(-) diff --git a/pkg/monitor/tsdb/driver/victoriametrics/result_parser.go b/pkg/monitor/tsdb/driver/victoriametrics/result_parser.go index 95a9eb1db0..e199162fe0 100644 --- a/pkg/monitor/tsdb/driver/victoriametrics/result_parser.go +++ b/pkg/monitor/tsdb/driver/victoriametrics/result_parser.go @@ -58,6 +58,7 @@ func (p *points) add(op *points) error { return errors.Errorf("input values' are %#v, which length isn't equal 2", oVal) } val = append(val, oVal[1]) + p.values[i] = val } return nil } @@ -69,7 +70,7 @@ func (p *points) isEqual(op *points) bool { return reflect.DeepEqual(p.columns, op.columns) && reflect.DeepEqual(p.tags, op.tags) && reflect.DeepEqual(p.values, op.values) } -func newPointsByResult(result ResponseDataResult) (*points, error) { +func newPointsByResult(result ResponseDataResult, sameTimes sets.String) (*points, error) { tags := result.Metric column, ok := tags[translator.UNION_RESULT_NAME] if !ok { @@ -83,10 +84,18 @@ func newPointsByResult(result ResponseDataResult) (*points, error) { } values := result.Values id := newMapId(tags) + filterValues := []ResponseDataResultValue{} + for _, val := range values { + valTime := fmt.Sprintf("%s", val[0]) + if sameTimes.Has(valTime) { + tmpVal := val + filterValues = append(filterValues, tmpVal) + } + } return &points{ id: id, columns: []string{column}, - values: values, + values: filterValues, tags: tags, }, nil } @@ -95,8 +104,22 @@ func newPointsByResults(results []ResponseDataResult) ([]*points, error) { uniq := make(map[string]*points, 0) ret := make([]*points, 0) + var sameTimes sets.String = nil for _, result := range results { - p, err := newPointsByResult(result) + resultTime := sets.NewString() + for _, v := range result.Values { + cTime := fmt.Sprintf("%v", v[0]) + resultTime.Insert(cTime) + } + if sameTimes == nil { + sameTimes = resultTime + } else { + sameTimes = sameTimes.Intersection(resultTime) + } + } + + for _, result := range results { + p, err := newPointsByResult(result, sameTimes) if err != nil { return nil, errors.Wrapf(err, "new points by result: %#v", result) } diff --git a/pkg/monitor/tsdb/driver/victoriametrics/vm.go b/pkg/monitor/tsdb/driver/victoriametrics/vm.go index 7ca7ca1d90..a4a334458c 100644 --- a/pkg/monitor/tsdb/driver/victoriametrics/vm.go +++ b/pkg/monitor/tsdb/driver/victoriametrics/vm.go @@ -192,8 +192,8 @@ func parseTimepoint(val ResponseDataResultValue) (tsdb.TimePoint, error) { valStr := val[i] pVal := parsePointValue(valStr) timepoint = append(timepoint, pVal) - timepoint = append(timepoint, timestamp) } + timepoint = append(timepoint, timestamp) return timepoint, nil }