Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 22 additions & 11 deletions bigtable/metrics_monitoring_exporter.go
Original file line number Diff line number Diff line change
Expand Up @@ -249,11 +249,15 @@ func (me *monitoringExporter) recordToTimeSeriesPb(m otelmetricdata.Metrics) ([]
tss = append(tss, ts)
}
case otelmetricdata.Sum[int64]:
kind := googlemetricpb.MetricDescriptor_CUMULATIVE
if !a.IsMonotonic {
kind = googlemetricpb.MetricDescriptor_GAUGE
}
for _, point := range a.DataPoints {
metric, mr := me.recordToMetricAndMonitoredResourcePbs(m, point.Attributes)
var ts *monitoringpb.TimeSeries
var err error
ts, err = sumToTimeSeries[int64](point, m, mr)
ts, err = sumToTimeSeries[int64](point, m, mr, kind)
if err != nil {
errs = append(errs, err)
continue
Expand All @@ -267,16 +271,16 @@ func (me *monitoringExporter) recordToTimeSeriesPb(m otelmetricdata.Metrics) ([]
return tss, errors.Join(errs...)
}

func sumToTimeSeries[N int64 | float64](point otelmetricdata.DataPoint[N], metrics otelmetricdata.Metrics, mr *monitoredrespb.MonitoredResource) (*monitoringpb.TimeSeries, error) {
interval, err := toNonemptyTimeIntervalpb(point.StartTime, point.Time)
func sumToTimeSeries[N int64 | float64](point otelmetricdata.DataPoint[N], metrics otelmetricdata.Metrics, mr *monitoredrespb.MonitoredResource, kind googlemetricpb.MetricDescriptor_MetricKind) (*monitoringpb.TimeSeries, error) {
interval, err := toTimeIntervalPb(point.StartTime, point.Time, kind)
if err != nil {
return nil, err
}
value, valueType := numberDataPointToValue[N](point)
return &monitoringpb.TimeSeries{
Resource: mr,
Unit: string(metrics.Unit),
MetricKind: googlemetricpb.MetricDescriptor_CUMULATIVE,
MetricKind: kind,
ValueType: valueType,
Points: []*monitoringpb.Point{{
Interval: interval,
Expand All @@ -286,7 +290,7 @@ func sumToTimeSeries[N int64 | float64](point otelmetricdata.DataPoint[N], metri
}

func histogramToTimeSeries[N int64 | float64](point otelmetricdata.HistogramDataPoint[N], metrics otelmetricdata.Metrics, mr *monitoredrespb.MonitoredResource) (*monitoringpb.TimeSeries, error) {
interval, err := toNonemptyTimeIntervalpb(point.StartTime, point.Time)
interval, err := toTimeIntervalPb(point.StartTime, point.Time, googlemetricpb.MetricDescriptor_CUMULATIVE)
if err != nil {
return nil, err
}
Expand All @@ -307,13 +311,20 @@ func histogramToTimeSeries[N int64 | float64](point otelmetricdata.HistogramData
}, nil
}

func toNonemptyTimeIntervalpb(start, end time.Time) (*monitoringpb.TimeInterval, error) {
// The end time of a new interval must be at least a millisecond after the end time of the
// previous interval, for all non-gauge types.
// https://cloud.google.com/monitoring/api/ref_v3/rpc/google.monitoring.v3#timeinterval
if end.Sub(start).Milliseconds() <= 1 {
end = start.Add(time.Millisecond)
func toTimeIntervalPb(start, end time.Time, kind googlemetricpb.MetricDescriptor_MetricKind) (*monitoringpb.TimeInterval, error) {
if kind == googlemetricpb.MetricDescriptor_GAUGE {
// Same as https://github.com/googleapis/java-bigtable/pull/2719
start = end
} else { // CUMULATIVE or other types
// The end time of a new interval must be at least a millisecond after the end time of the
// previous interval for CUMULATIVE types.
// https://cloud.google.com/monitoring/api/ref_v3/rpc/google.monitoring.v3#timeinterval

if end.Sub(start) < time.Millisecond {
end = start.Add(time.Millisecond)
}
}

startpb := timestamppb.New(start)
endpb := timestamppb.New(end)
err := errors.Join(
Expand Down
23 changes: 21 additions & 2 deletions bigtable/metrics_monitoring_exporter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -430,7 +430,7 @@ func TestRecordToMpb(t *testing.T) {
func TestTimeIntervalStaggering(t *testing.T) {
var tm time.Time

interval, err := toNonemptyTimeIntervalpb(tm, tm)
interval, err := toTimeIntervalPb(tm, tm, googlemetricpb.MetricDescriptor_CUMULATIVE)
if err != nil {
t.Fatalf("conversion to PB failed: %v", err)
}
Expand All @@ -450,10 +450,29 @@ func TestTimeIntervalStaggering(t *testing.T) {
}
}

func TestTimeIntervalGauge(t *testing.T) {
startTime := time.Now()
endTime := startTime.Add(time.Second)

interval, err := toTimeIntervalPb(startTime, endTime, googlemetricpb.MetricDescriptor_GAUGE)
if err != nil {
t.Fatalf("conversion to PB failed: %v", err)
}

start := interval.StartTime.AsTime()
end := interval.EndTime.AsTime()

if !start.Equal(end) {
t.Errorf("Expected StartTime == EndTime for GAUGE, got StartTime=%v, EndTime=%v", start, end)
}
if !start.Equal(endTime) {
t.Errorf("Expected StartTime to be reset to EndTime for GAUGE, got StartTime=%v, expected %v", start, endTime)
}
}
func TestTimeIntervalPassthru(t *testing.T) {
var tm time.Time

interval, err := toNonemptyTimeIntervalpb(tm, tm.Add(time.Second))
interval, err := toTimeIntervalPb(tm, tm.Add(time.Second), googlemetricpb.MetricDescriptor_CUMULATIVE)
if err != nil {
t.Fatalf("conversion to PB failed: %v", err)
}
Expand Down
Loading