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
53 changes: 48 additions & 5 deletions sink/datadog_api.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"fmt"
"io"
"net/http"
"sort"
"strings"
"time"

Expand Down Expand Up @@ -182,16 +183,14 @@ func prepareDatadogPayload(ctx context.Context, series []datadogV2.MetricSeries,
}

var raw bytes.Buffer
grow := rawTarget
if estimated := maxSeries * 512; estimated < grow {
grow = estimated
}
candidates := min(maxSeries, len(series)-start)
grow := datadogRawCapacity(series[start:start+candidates], rawTarget)
if grow > 0 {
raw.Grow(grow)
}
raw.Write(datadogPayloadPrefix)

offsets := make([]int, 0, maxSeries)
offsets := make([]int, 0, candidates)
end := start
for end < len(series) && len(offsets) < maxSeries {
select {
Expand Down Expand Up @@ -282,6 +281,50 @@ func prepareDatadogPayload(ctx context.Context, series []datadogV2.MetricSeries,
return prepared, end, nil
}

// datadogRawCapacity uses the median of at most eight evenly spaced series.
// It inspects at most sixteen tags per series without scanning string contents.
// Allow extra room for JSON strings, but keep the existing small-series floor.
// This is only a hint: escaping, other fields, and unsampled outliers can still
// require growth. Exact payload checks remain authoritative.
func datadogRawCapacity(series []datadogV2.MetricSeries, target int) int {
count := len(series)
if count == 0 || target <= 0 {
return 0
}
if count > target/512 {
return target
}
samples := min(count, 8)
perSeries := target / count
// perSeries >= 512. Cap before doubling or multiplying to avoid overflow.
stringBudget := (perSeries - 128) / 2
var sizes [8]int
for i := 0; i < samples; i++ {
// Include both ends; quotient/remainder avoids multiplying count by i.
divisor := max(samples-1, 1)
index := (count-1)/divisor*i + (count-1)%divisor*i/divisor
value := &series[index]
size := min(len(value.Metric), stringBudget)
tags := min(len(value.Tags), 16)
for j := 0; j < tags; j++ {
divisor := max(tags-1, 1)
index := (len(value.Tags)-1)/divisor*j + (len(value.Tags)-1)%divisor*j/divisor
size += min(len(value.Tags[index]), stringBudget-size)
}
sizes[i] = max(512, 128+2*size)
if size == stringBudget {
sizes[i] = perSeries
}
}
sort.Ints(sizes[:samples])
// The upper median avoids extrapolating a single unusually large series.
estimate := sizes[samples/2]
if estimate == perSeries {
return target
}
return estimate * count
}

func payloadPrefixForSeries(raw []byte, end int) []byte {
prefix := make([]byte, end+len(datadogPayloadSuffix))
copy(prefix, raw[:end])
Expand Down
96 changes: 96 additions & 0 deletions sink/datadog_capacity_hint_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
package sink

import (
"context"
"encoding/json"
"fmt"
"strings"
"testing"

"github.com/DataDog/datadog-api-client-go/v2/api/datadogV2"
"github.com/stretchr/testify/require"
)

func TestDatadogRawCapacityBounds(t *testing.T) {
for _, n := range []int{0, 1, 2, 8, 9, 61, 383, 1000, 10000} {
series := testMetricSeries(n)
for _, target := range []int{0, 1, 511, 512, 513, 4095, 4718592, int(^uint(0) >> 1)} {
capacity := datadogRawCapacity(series, target)
require.Equal(t, min(n*512, target), capacity, "ordinary series retain the existing hint")
if n == 0 {
continue
}
series[n-1].Tags = []string{strings.Repeat("large", 20000)}
capacity = datadogRawCapacity(series, target)
require.GreaterOrEqual(t, capacity, 0)
require.LessOrEqual(t, capacity, target)
series[n-1] = testMetricSeries(1)[0]
}
}
}

func TestDatadogRawCapacityIsolatedOutlier(t *testing.T) {
for _, position := range []int{0, 499, 999} {
series := testMetricSeries(1000)
series[position].Tags = []string{strings.Repeat("x", 65536)}
require.Equal(t, 512000, datadogRawCapacity(series, datadogTargetDecompressedPayloadSize))
}
}

func TestDatadogRawCapacityFindsLateTags(t *testing.T) {
series := testMetricSeries(1000)
for i := range series {
series[i].Tags = make([]string, 64)
series[i].Tags[63] = strings.Repeat("x", 4096)
}
require.Equal(t, datadogTargetDecompressedPayloadSize, datadogRawCapacity(series, datadogTargetDecompressedPayloadSize))
}

func TestDatadogCapacityHardLimitSeries(t *testing.T) {
for _, compress := range []bool{false, true} {
limit := datadogMaxCompressedPayloadSize
if compress {
limit = datadogMaxDecompressedPayloadSize
}
for _, delta := range []int{-1, 0, 1} {
t.Run(fmt.Sprintf("gzip=%t/delta=%d", compress, delta), func(t *testing.T) {
series := testMetricSeries(1)
series[0].Tags = []string{""}
encoded, err := json.Marshal(series[0])
require.NoError(t, err)
series[0].Tags[0] = strings.Repeat("x", limit+delta-len(encoded)-len(datadogPayloadPrefix)-len(datadogPayloadSuffix))
payload, _, err := prepareDatadogPayload(context.Background(), series, 0, 10000, compress, defaultDatadogPayloadLimits())
if delta > 0 {
require.Error(t, err)
return
}
require.NoError(t, err)
require.Equal(t, limit+delta, payload.uncompressedBytes)
decoded, err := decodePreparedMetricPayload(payload)
require.NoError(t, err)
require.Equal(t, series, decoded.Series)
})
}
}
}

func BenchmarkDatadogRawCapacityHint(b *testing.B) {
for _, n := range []int{61, 383, 1000, 10000} {
for _, tags := range []int{2, 64} {
b.Run(fmt.Sprintf("%d/tags=%d", n, tags), func(b *testing.B) {
series := make([]datadogV2.MetricSeries, n)
for i := range series {
series[i].Tags = make([]string, tags)
for j := range series[i].Tags {
series[i].Tags[j] = "tag:synthetic"
}
}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
datadogRawCapacity(series, datadogTargetDecompressedPayloadSize)
}
})
}
}
}
Loading
Loading