From cdbe35c185bfaace45dc4bccee1b021074d22d87 Mon Sep 17 00:00:00 2001 From: Qiheng He Date: Fri, 21 Aug 2026 15:12:27 +0800 Subject: [PATCH] Add support for InfluxDB 2 and InfluxDB 3 Core --- cmd/go.mod | 9 +- cmd/go.sum | 13 ++ cmd/internal/storage/influxdb/influxdb.go | 147 +++++++++---- .../storage/influxdb/influxdb_test.go | 36 +++- .../storage/influxdb/influxdb_token_test.go | 195 ++++++++++++++++++ docs/storage/influxdb.md | 74 ++++++- 6 files changed, 421 insertions(+), 53 deletions(-) create mode 100644 cmd/internal/storage/influxdb/influxdb_token_test.go diff --git a/cmd/go.mod b/cmd/go.mod index 901b5569c5..431962d414 100644 --- a/cmd/go.mod +++ b/cmd/go.mod @@ -119,6 +119,13 @@ require ( gopkg.in/yaml.v3 v3.0.1 // indirect ) -require github.com/euank/go-kmsg-parser v2.0.0+incompatible // indirect +require github.com/influxdata/influxdb-client-go/v2 v2.14.0 + +require ( + github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect + github.com/euank/go-kmsg-parser v2.0.0+incompatible // indirect + github.com/influxdata/line-protocol v0.0.0-20200327222509-2487e7298839 // indirect + github.com/oapi-codegen/runtime v1.0.0 // indirect +) replace github.com/google/cadvisor/lib => ../lib diff --git a/cmd/go.sum b/cmd/go.sum index 9c58875475..51af191dcc 100644 --- a/cmd/go.sum +++ b/cmd/go.sum @@ -6,6 +6,7 @@ cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdB cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10= github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk= github.com/SeanDolphin/bqschema v1.0.0 h1:iCYFd5Qsw6caM2k5/SsITSL9+3kQCr+oz6pnNjWTq90= github.com/SeanDolphin/bqschema v1.0.0/go.mod h1:TYInVncsPIZH7kybQoIUNJ4pFX1cUc8LoP9RSOxIs6c= github.com/Shopify/sarama v1.38.1 h1:lqqPUPQZ7zPqYlWpTh+LQ9bhYNu2xJL6k1SJN4WVe2A= @@ -18,6 +19,8 @@ github.com/andreyvit/diff v0.0.0-20170406064948-c7f18ee00883 h1:bvNMNQO63//z+xNg github.com/andreyvit/diff v0.0.0-20170406064948-c7f18ee00883/go.mod h1:rCTlJbsFo29Kk6CurOXKm700vrz8f0KW0JNfpkRJY/8= github.com/apache/arrow/go/v7 v7.0.1 h1:WpCfq+AQxvXaI6/KplHE27MPMFx5av0o5NbPCTAGfy4= github.com/apache/arrow/go/v7 v7.0.1/go.mod h1:JxDpochJbCVxqbX4G8i1jRqMrnTCQdf8pTccAfLD8Es= +github.com/apapsch/go-jsonmerge/v2 v2.0.0 h1:axGnT1gRIfimI7gJifB699GoE/oq+F2MU7Dml6nw9rQ= +github.com/apapsch/go-jsonmerge/v2 v2.0.0/go.mod h1:lvDnEdqiQrp0O42VQGgmlKpxL1AP2+08jFMw88y4klk= github.com/aws/aws-sdk-go-v2 v1.36.3 h1:mJoei2CxPutQVxaATCzDUjcZEjVRdpsiiXi2o38yqWM= github.com/aws/aws-sdk-go-v2 v1.36.3/go.mod h1:LLXuLpgzEbD766Z5ECcRmi8AzSwfZItDtmABVkRLGzg= github.com/aws/aws-sdk-go-v2/config v1.29.14 h1:f+eEi/2cKCg9pqKBoAIwRGzVb70MRKqWX4dg1BDcSJM= @@ -50,6 +53,7 @@ github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= github.com/blang/semver/v4 v4.0.0/go.mod h1:IbckMUScFkM3pff0VJDNKRiT6TG/YpiHIM2yvyW5YoQ= +github.com/bmatcuk/doublestar v1.1.1/go.mod h1:UD6OnuiIn0yFxxA2le/rnRU1G4RaI4UvFv1sNto9p6w= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/containerd/containerd/api v1.10.0 h1:5n0oHYVBwN4VhoX9fFykCV9dF1/BvAXeg2F8W6UYq1o= @@ -147,8 +151,12 @@ github.com/influxdata/flux v0.196.1 h1:RZypfrrAZeIixD/2As+qoyrwfs6Gx7VcnTRNGNjBf github.com/influxdata/flux v0.196.1/go.mod h1:+Y4mBygx6q98onpdKJd6vJPrTNjHriQhwh/gM+3IvUQ= github.com/influxdata/influxdb v1.12.0 h1:hVFeHtEUh/MkI6YCmB1TAXlPgV4ZNh/V5FP6z7ywRLo= github.com/influxdata/influxdb v1.12.0/go.mod h1:11RjLuBNkuWaJQFViRF/rpNzICfU6X0nuO003yeleKY= +github.com/influxdata/influxdb-client-go/v2 v2.14.0 h1:AjbBfJuq+QoaXNcrova8smSjwJdUHnwvfjMF71M1iI4= +github.com/influxdata/influxdb-client-go/v2 v2.14.0/go.mod h1:Ahpm3QXKMJslpXl3IftVLVezreAUtBOTZssDrjZEFHI= github.com/influxdata/influxql v1.4.1 h1:UB+TMc9cB6mDdkPmH/5sBU8FQ+ZCWRX2JPcPDIFrLcs= github.com/influxdata/influxql v1.4.1/go.mod h1:VqxAKyQz5p8GzgGsxWalCWYGxEqw6kvJo2IickMQiQk= +github.com/influxdata/line-protocol v0.0.0-20200327222509-2487e7298839 h1:W9WBk7wlPfJLvMCdtV4zPulc4uCPrlywQOmbFOhgQNU= +github.com/influxdata/line-protocol v0.0.0-20200327222509-2487e7298839/go.mod h1:xaLFMmpvUxqXtVkUJfg9QmT88cDaCJ3ZKgdZ78oO8Qo= github.com/influxdb/influxdb v1.7.9 h1:KMBwwvyJyBppIwrg5t0662p+Yei/ucnIkqUl8txiQdQ= github.com/influxdb/influxdb v1.7.9/go.mod h1:GpjLgHRqWhDGlPAg7+Rj6NAYuzPojBM8XLG5Ouvvq+Q= github.com/jcmturner/aescts/v2 v2.0.0 h1:9YKLH6ey7H4eDBXW8khjYslgyqG2xZikXP0EQFKrle8= @@ -163,6 +171,7 @@ github.com/jcmturner/gokrb5/v8 v8.4.4 h1:x1Sv4HaTpepFkXbt2IkL29DXRf8sOfZXo8eRKh6 github.com/jcmturner/gokrb5/v8 v8.4.4/go.mod h1:1btQEpgT6k+unzCwX1KdWMEwPPkkgBtP+F6aCACiMrs= github.com/jcmturner/rpc/v2 v2.0.3 h1:7FXXj8Ti1IaVFpSAziCZWNzbNuZmnvw/i6CqLNdWfZY= github.com/jcmturner/rpc/v2 v2.0.3/go.mod h1:VUJYCIDm3PVOEHw8sgt091/20OJjskO/YJki3ELg/Hc= +github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPciBPrBUjwbNvtwB6RQlve+hkpll6QSNmOE= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -188,6 +197,8 @@ github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8m github.com/nxadm/tail v1.4.4/go.mod h1:kenIhsEOeOJmVchQTgglprH7qJGnHDVpk1VPCcaMI8A= github.com/nxadm/tail v1.4.8 h1:nPr65rt6Y5JFSKQO7qToXr7pePgD6Gwiw05lkbyAQTE= github.com/nxadm/tail v1.4.8/go.mod h1:+ncqLTQzXmGhMZNUePPaPqPvBxHAIsmXswZKocGu+AU= +github.com/oapi-codegen/runtime v1.0.0 h1:P4rqFX5fMFWqRzY9M/3YF9+aPSPPB06IzP2P7oOxrWo= +github.com/oapi-codegen/runtime v1.0.0/go.mod h1:LmCUMQuPB4M/nLXilQXhHw+BLZdDb18B34OO356yJ/A= github.com/onsi/ginkgo v1.6.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.12.1/go.mod h1:zj2OWP4+oCPe1qIXoGWkgMRwljMUYCdkwsT2108oapk= github.com/onsi/ginkgo v1.16.5 h1:8xi0RTUf59SOSfEtZMvwTvXYMzG4gV23XVHOZiXNtnE= @@ -230,11 +241,13 @@ github.com/sergi/go-diff v1.1.0 h1:we8PVUC3FE2uYfodKH/nBHMSetSfHDR6scGdBi+erh0= github.com/sergi/go-diff v1.1.0/go.mod h1:STckp+ISIX8hZLjrqAeVduY0gWCT9IjLuqbuNXdaHfM= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= +github.com/spkg/bom v0.0.0-20160624110644-59b7046e48ad/go.mod h1:qLr4V1qq6nMqFKkMo8ZTx3f+BZEkzsRUY10Xsm2mwU0= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= github.com/stretchr/testify v1.5.1/go.mod h1:5W2xD1RspED5o8YsWQXVCued0rvSQ+mT+I5cxcmMvtA= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= diff --git a/cmd/internal/storage/influxdb/influxdb.go b/cmd/internal/storage/influxdb/influxdb.go index 793512576b..694f11e565 100644 --- a/cmd/internal/storage/influxdb/influxdb.go +++ b/cmd/internal/storage/influxdb/influxdb.go @@ -15,6 +15,7 @@ package influxdb import ( + "context" "flag" "fmt" "net/url" @@ -26,6 +27,9 @@ import ( "github.com/google/cadvisor/lib/storage" "github.com/google/cadvisor/lib/version" + influxdb2 "github.com/influxdata/influxdb-client-go/v2" + api "github.com/influxdata/influxdb-client-go/v2/api" + "github.com/influxdata/influxdb-client-go/v2/api/write" influxdb "github.com/influxdb/influxdb/client" ) @@ -33,12 +37,20 @@ func init() { storage.RegisterStorageDriver("influxdb", new) } -var argDbRetentionPolicy = flag.String("storage_driver_influxdb_retention_policy", "", "retention policy") +var ( + argDbRetentionPolicy = flag.String("storage_driver_influxdb_retention_policy", "", "retention policy") + argDbAuthToken = flag.String("storage_driver_influxdb_auth_token", "", "InfluxDB API token. When set, stats are written through the InfluxDB 2.x API, which is also implemented by InfluxDB 3 Core and InfluxDB Cloud, using token authentication instead of the InfluxDB 1.x username/password API") + argDbOrg = flag.String("storage_driver_influxdb_org", "", "InfluxDB organization to write to. Required by InfluxDB 2.x, ignored by InfluxDB 3 Core") + argDbBucket = flag.String("storage_driver_influxdb_bucket", "", "InfluxDB 2.x bucket or InfluxDB 3 Core database to write to. Defaults to the -storage_driver_db value") +) type influxdbStorage struct { client *influxdb.Client + v2Client influxdb2.Client + writeAPI api.WriteAPIBlocking machineName string database string + bucket string retentionPolicy string bufferDuration time.Duration lastWrite time.Time @@ -120,6 +132,9 @@ func new() (storage.StorageDriver, error) { *argDbRetentionPolicy, *storage.ArgDbUsername, *storage.ArgDbPassword, + *argDbAuthToken, + *argDbOrg, + *argDbBucket, *storage.ArgDbHost, *storage.ArgDbIsSecure, *storage.ArgDbBufferDuration, @@ -208,31 +223,35 @@ func (s *influxdbStorage) containerStatsToPoints( stats *info.ContainerStats, ) (points []*influxdb.Point) { // CPU usage: Total usage in nanoseconds - points = append(points, makePoint(serCPUUsageTotal, stats.Cpu.Usage.Total)) + if stats.Cpu != nil { + points = append(points, makePoint(serCPUUsageTotal, stats.Cpu.Usage.Total)) - // CPU usage: Time spend in system space (in nanoseconds) - points = append(points, makePoint(serCPUUsageSystem, stats.Cpu.Usage.System)) + // CPU usage: Time spend in system space (in nanoseconds) + points = append(points, makePoint(serCPUUsageSystem, stats.Cpu.Usage.System)) - // CPU usage: Time spent in user space (in nanoseconds) - points = append(points, makePoint(serCPUUsageUser, stats.Cpu.Usage.User)) + // CPU usage: Time spent in user space (in nanoseconds) + points = append(points, makePoint(serCPUUsageUser, stats.Cpu.Usage.User)) - // CPU usage per CPU - for i := 0; i < len(stats.Cpu.Usage.PerCpu); i++ { - point := makePoint(serCPUUsagePerCPU, stats.Cpu.Usage.PerCpu[i]) - tags := map[string]string{"instance": fmt.Sprintf("%v", i)} - addTagsToPoint(point, tags) + // CPU usage per CPU + for i := 0; i < len(stats.Cpu.Usage.PerCpu); i++ { + point := makePoint(serCPUUsagePerCPU, stats.Cpu.Usage.PerCpu[i]) + tags := map[string]string{"instance": fmt.Sprintf("%v", i)} + addTagsToPoint(point, tags) - points = append(points, point) - } + points = append(points, point) + } - // Load Average - points = append(points, makePoint(serLoadAverage, stats.Cpu.LoadAverage)) + // Load Average + points = append(points, makePoint(serLoadAverage, stats.Cpu.LoadAverage)) + } // Network Stats - points = append(points, makePoint(serRxBytes, stats.Network.RxBytes)) - points = append(points, makePoint(serRxErrors, stats.Network.RxErrors)) - points = append(points, makePoint(serTxBytes, stats.Network.TxBytes)) - points = append(points, makePoint(serTxErrors, stats.Network.TxErrors)) + if stats.Network != nil { + points = append(points, makePoint(serRxBytes, stats.Network.RxBytes)) + points = append(points, makePoint(serRxErrors, stats.Network.RxErrors)) + points = append(points, makePoint(serTxBytes, stats.Network.TxBytes)) + points = append(points, makePoint(serTxErrors, stats.Network.TxErrors)) + } // Referenced Memory points = append(points, makePoint(serReferencedMemory, stats.ReferencedMemory)) @@ -246,6 +265,9 @@ func (s *influxdbStorage) memoryStatsToPoints( cInfo *info.ContainerInfo, stats *info.ContainerStats, ) (points []*influxdb.Point) { + if stats.Memory == nil { + return points + } // Memory Usage points = append(points, makePoint(serMemoryUsage, stats.Memory.Usage)) // Maximum memory usage recorded @@ -414,35 +436,63 @@ func (s *influxdbStorage) AddStats(cInfo *info.ContainerInfo, stats *info.Contai } }() if len(pointsToFlush) > 0 { - points := make([]influxdb.Point, len(pointsToFlush)) - for i, p := range pointsToFlush { - points[i] = *p + if err := s.flush(pointsToFlush, stats.Timestamp); err != nil { + return err } + } + return nil +} - batchTags := map[string]string{tagMachineName: s.machineName} - bp := influxdb.BatchPoints{ - Points: points, - Database: s.database, - RetentionPolicy: s.retentionPolicy, - Tags: batchTags, - Time: stats.Timestamp, +func (s *influxdbStorage) flush(pointsToFlush []*influxdb.Point, timestamp time.Time) error { + if s.writeAPI != nil { + // Token-authenticated mode: write through the InfluxDB 2.x API, + // which is also implemented by InfluxDB 3 Core. + points := make([]*write.Point, len(pointsToFlush)) + for i, p := range pointsToFlush { + points[i] = influxdb2.NewPoint(p.Measurement, p.Tags, p.Fields, p.Time) } - response, err := s.client.Write(bp) - if err != nil || checkResponseForErrors(response) != nil { + if err := s.writeAPI.WritePoint(context.Background(), points...); err != nil { return fmt.Errorf("failed to write stats to influxDb - %s", err) } + return nil + } + + points := make([]influxdb.Point, len(pointsToFlush)) + for i, p := range pointsToFlush { + points[i] = *p + } + + batchTags := map[string]string{tagMachineName: s.machineName} + bp := influxdb.BatchPoints{ + Points: points, + Database: s.database, + RetentionPolicy: s.retentionPolicy, + Tags: batchTags, + Time: timestamp, + } + response, err := s.client.Write(bp) + if err != nil || checkResponseForErrors(response) != nil { + return fmt.Errorf("failed to write stats to influxDb - %s", err) } return nil } func (s *influxdbStorage) Close() error { + if s.v2Client != nil { + s.v2Client.Close() + } s.client = nil + s.v2Client = nil + s.writeAPI = nil return nil } // machineName: A unique identifier to identify the host that current cAdvisor // instance is running on. // influxdbHost: The host which runs influxdb (host:port) +// authToken: When non-empty, stats are written with token authentication +// through the InfluxDB 2.x API (also implemented by InfluxDB 3 Core) and +// username/password are ignored. The bucket falls back to database when empty. func newStorage( machineName, tablename, @@ -450,6 +500,9 @@ func newStorage( retentionPolicy, username, password, + authToken, + org, + bucket, influxdbHost string, isSecure bool, bufferDuration time.Duration, @@ -462,19 +515,7 @@ func newStorage( url.Scheme = "https" } - config := &influxdb.Config{ - URL: *url, - Username: username, - Password: password, - UserAgent: fmt.Sprintf("%v/%v", "cAdvisor", version.Info["version"]), - } - client, err := influxdb.NewClient(*config) - if err != nil { - return nil, err - } - ret := &influxdbStorage{ - client: client, machineName: machineName, database: database, retentionPolicy: retentionPolicy, @@ -482,6 +523,26 @@ func newStorage( lastWrite: time.Now(), points: make([]*influxdb.Point, 0), } + if authToken != "" { + if bucket == "" { + bucket = database + } + ret.bucket = bucket + ret.v2Client = influxdb2.NewClient(url.String(), authToken) + ret.writeAPI = ret.v2Client.WriteAPIBlocking(org, bucket) + } else { + config := &influxdb.Config{ + URL: *url, + Username: username, + Password: password, + UserAgent: fmt.Sprintf("%v/%v", "cAdvisor", version.Info["version"]), + } + client, err := influxdb.NewClient(*config) + if err != nil { + return nil, err + } + ret.client = client + } ret.readyToFlush = ret.defaultReadyToFlush return ret, nil } diff --git a/cmd/internal/storage/influxdb/influxdb_test.go b/cmd/internal/storage/influxdb/influxdb_test.go index 1a9cf7e7ea..cfc7fedda9 100644 --- a/cmd/internal/storage/influxdb/influxdb_test.go +++ b/cmd/internal/storage/influxdb/influxdb_test.go @@ -126,6 +126,9 @@ func runStorageTest(f func(test.TestStorageDriver, *testing.T), t *testing.T, bu retentionPolicy, username, password, + "", + "", + "", hostname, false, time.Duration(bufferCount)) @@ -147,6 +150,9 @@ func runStorageTest(f func(test.TestStorageDriver, *testing.T), t *testing.T, bu retentionPolicy, username, password, + "", + "", + "", hostname, false, time.Duration(bufferCount)) @@ -204,6 +210,9 @@ func TestContainerFileSystemStatsToPoints(t *testing.T) { retentionPolicy, username, password, + "", + "", + "", influxdbHost, false, 2*time.Minute) assert.Nil(err) @@ -240,15 +249,15 @@ func TestContainerStatsToPoints(t *testing.T) { // Then assert.NotEmpty(t, points) - assert.Len(t, points, 34+len(stats.Cpu.Usage.PerCpu)) + assert.Len(t, points, 36+len(stats.Cpu.Usage.PerCpu)) // CPU stats - assertContainsPointWithValue(t, points, serCpuUsageTotal, stats.Cpu.Usage.Total) - assertContainsPointWithValue(t, points, serCpuUsageSystem, stats.Cpu.Usage.System) - assertContainsPointWithValue(t, points, serCpuUsageUser, stats.Cpu.Usage.User) + assertContainsPointWithValue(t, points, serCPUUsageTotal, stats.Cpu.Usage.Total) + assertContainsPointWithValue(t, points, serCPUUsageSystem, stats.Cpu.Usage.System) + assertContainsPointWithValue(t, points, serCPUUsageUser, stats.Cpu.Usage.User) assertContainsPointWithValue(t, points, serLoadAverage, stats.Cpu.LoadAverage) for _, cpu_usage := range stats.Cpu.Usage.PerCpu { - assertContainsPointWithValue(t, points, serCpuUsagePerCpu, cpu_usage) + assertContainsPointWithValue(t, points, serCPUUsagePerCPU, cpu_usage) } // Memory stats @@ -279,7 +288,7 @@ func TestContainerStatsToPoints(t *testing.T) { assertContainsPointWithValue(t, points, serRxBytes, stats.Network.RxBytes) assertContainsPointWithValue(t, points, serRxErrors, stats.Network.RxErrors) assertContainsPointWithValue(t, points, serTxBytes, stats.Network.TxBytes) - assertContainsPointWithValue(t, points, serTxBytes, stats.Network.TxErrors) + assertContainsPointWithValue(t, points, serTxErrors, stats.Network.TxErrors) // Perf stats for _, perfStat := range stats.PerfStats { @@ -327,6 +336,9 @@ func createTestStorage() (*influxdbStorage, error) { retentionPolicy, username, password, + "", + "", + "", influxdbHost, false, 2*time.Minute) @@ -350,11 +362,11 @@ func createTestStats() (*info.ContainerInfo, *info.ContainerStats) { stats := &info.ContainerStats{ Timestamp: time.Now(), - Cpu: info.CpuStats{ + Cpu: &info.CpuStats{ Usage: cpuUsage, LoadAverage: int32(rand.Intn(1000)), }, - Memory: info.MemoryStats{ + Memory: &info.MemoryStats{ Usage: 26767396864, MaxUsage: 30429605888, Cache: 7837376512, @@ -372,6 +384,14 @@ func createTestStats() (*info.ContainerInfo, *info.ContainerStats) { "1GB": {Usage: 1234, MaxUsage: 5678, Failcnt: 9}, "2GB": {Usage: 9876, MaxUsage: 5432, Failcnt: 1}, }, + Network: &info.NetworkStats{ + InterfaceStats: info.InterfaceStats{ + RxBytes: 123, + RxErrors: 4, + TxBytes: 456, + TxErrors: 7, + }, + }, ReferencedMemory: 12345, PerfStats: []info.PerfStat{{Cpu: 1, PerfValue: info.PerfValue{Name: "cycles", ScalingRatio: 1.5, Value: 4589}}}, Resctrl: info.ResctrlStats{ diff --git a/cmd/internal/storage/influxdb/influxdb_token_test.go b/cmd/internal/storage/influxdb/influxdb_token_test.go new file mode 100644 index 0000000000..6d57d0c83e --- /dev/null +++ b/cmd/internal/storage/influxdb/influxdb_token_test.go @@ -0,0 +1,195 @@ +// Copyright 2026 Google Inc. All Rights Reserved. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package influxdb + +import ( + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + info "github.com/google/cadvisor/info/v1" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Tests for the token-authenticated InfluxDB 2.x / InfluxDB 3 Core write +// path. Unlike influxdb_test.go (build tag influxdb_test), these tests do not +// require a running InfluxDB instance and run as part of the default +// `go test ./...` invocation. + +func tokenModeTestStats() (*info.ContainerInfo, *info.ContainerStats) { + cInfo := &info.ContainerInfo{ + ContainerReference: info.ContainerReference{ + Name: "testContainername", + Aliases: []string{"testContainerAlias"}, + }, + } + stats := &info.ContainerStats{ + Timestamp: time.Unix(1700000000, 0).UTC(), + Cpu: &info.CpuStats{ + Usage: info.CpuUsage{Total: 100, User: 40, System: 60}, + }, + Memory: &info.MemoryStats{Usage: 2048, WorkingSet: 1024}, + Network: &info.NetworkStats{ + InterfaceStats: info.InterfaceStats{ + RxBytes: 11, + RxErrors: 1, + TxBytes: 22, + TxErrors: 2, + }, + }, + } + return cInfo, stats +} + +func TestNewStorageLegacyModeWithoutAuthToken(t *testing.T) { + storage, err := newStorage("machineA", + "cadvisor_table", + "cadvisor_test", + "cadvisor_test_rp", + "root", + "root", + "", + "", + "", + "localhost:8086", + false, + time.Minute) + require.NoError(t, err) + defer storage.Close() + + assert.NotNil(t, storage.client) + assert.Nil(t, storage.v2Client) + assert.Nil(t, storage.writeAPI) +} + +func TestNewStorageTokenMode(t *testing.T) { + storage, err := newStorage("machineA", + "cadvisor_table", + "cadvisor_test", + "cadvisor_test_rp", + "root", + "root", + "my-token", + "my-org", + "", + "localhost:8086", + true, + time.Minute) + require.NoError(t, err) + defer storage.Close() + + assert.Nil(t, storage.client) + assert.NotNil(t, storage.v2Client) + assert.NotNil(t, storage.writeAPI) + // The bucket falls back to the database name when unset. + assert.Equal(t, "cadvisor_test", storage.bucket) +} + +func TestNewStorageTokenModeExplicitBucket(t *testing.T) { + storage, err := newStorage("machineA", + "cadvisor_table", + "cadvisor_test", + "cadvisor_test_rp", + "root", + "root", + "my-token", + "my-org", + "my-bucket", + "localhost:8086", + false, + time.Minute) + require.NoError(t, err) + defer storage.Close() + + assert.Equal(t, "my-bucket", storage.bucket) +} + +func TestAddStatsWithAuthTokenWritesToV2API(t *testing.T) { + var gotPath, gotAuthHeader, gotOrg, gotBucket, gotBody string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + gotAuthHeader = r.Header.Get("Authorization") + gotOrg = r.URL.Query().Get("org") + gotBucket = r.URL.Query().Get("bucket") + body, err := io.ReadAll(r.Body) + require.NoError(t, err) + gotBody = string(body) + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + storage, err := newStorage("machineA", + "cadvisor_table", + "cadvisor_db", + "", + "ignored-user", + "ignored-password", + "my-token", + "my-org", + "", + strings.TrimPrefix(server.URL, "http://"), + false, + 0) + require.NoError(t, err) + defer storage.Close() + + cInfo, stats := tokenModeTestStats() + require.NoError(t, storage.AddStats(cInfo, stats)) + + assert.Equal(t, "/api/v2/write", gotPath) + assert.Equal(t, "Token my-token", gotAuthHeader) + assert.Equal(t, "my-org", gotOrg) + // The bucket falls back to the -storage_driver_db value. + assert.Equal(t, "cadvisor_db", gotBucket) + + // The batch is written as line protocol with the cAdvisor tags attached. + assert.Contains(t, gotBody, "cpu_usage_total") + assert.Contains(t, gotBody, "memory_usage") + assert.Contains(t, gotBody, "rx_bytes") + assert.Contains(t, gotBody, "machine=machineA") + assert.Contains(t, gotBody, "container_name=testContainerAlias") +} + +func TestAddStatsWithAuthTokenReturnsWriteErrors(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + http.Error(w, `{"code":"unauthorized","message":"Unauthorized"}`, http.StatusUnauthorized) + })) + defer server.Close() + + storage, err := newStorage("machineA", + "cadvisor_table", + "cadvisor_db", + "", + "ignored-user", + "ignored-password", + "invalid-token", + "my-org", + "", + strings.TrimPrefix(server.URL, "http://"), + false, + 0) + require.NoError(t, err) + defer storage.Close() + + cInfo, stats := tokenModeTestStats() + err = storage.AddStats(cInfo, stats) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to write stats to influxDb") +} diff --git a/docs/storage/influxdb.md b/docs/storage/influxdb.md index 1506eb451e..ea1ef7c49e 100644 --- a/docs/storage/influxdb.md +++ b/docs/storage/influxdb.md @@ -1,6 +1,6 @@ # Exporting cAdvisor Stats to InfluxDB -cAdvisor supports exporting stats to [InfluxDB](http://influxdb.com). To use InfluxDB, you need to pass some additional flags to cAdvisor telling it where the InfluxDB instance is located: +cAdvisor supports exporting stats to [InfluxDB](https://www.influxdata.com/) 1.x, 2.x, InfluxDB 3 Core and InfluxDB Cloud Serverless. Set the storage driver as InfluxDB. @@ -8,6 +8,8 @@ Set the storage driver as InfluxDB. -storage_driver=influxdb ``` +## InfluxDB 1.x + Specify what InfluxDB instance to push data to: ``` @@ -27,6 +29,76 @@ Specify what InfluxDB instance to push data to: -storage_driver_influxdb_retention_policy ``` +## InfluxDB 2.x, InfluxDB 3 Core and InfluxDB Cloud + +InfluxDB 2.x and later do not use usernames and passwords for API access; they use API tokens instead. To write through the InfluxDB v2 API (which is also implemented by InfluxDB 3 Core), set `-storage_driver_influxdb_auth_token`. The `-storage_driver_user` and `-storage_driver_password` flags are then ignored. + +``` + # The *ip:port* of the database. Default is 'localhost:8086' (InfluxDB 3 Core defaults to 'localhost:8181') + -storage_driver_host=ip:port + # database name. Used as the InfluxDB 2.x bucket (or the InfluxDB 3 Core database) unless -storage_driver_influxdb_bucket is set. Uses db 'cadvisor' by default + -storage_driver_db + # InfluxDB API token. Required to enable token-based writes. + -storage_driver_influxdb_auth_token + # InfluxDB organization. Required by InfluxDB 2.x; ignored by InfluxDB 3 Core. + -storage_driver_influxdb_org + # InfluxDB 2.x bucket or InfluxDB 3 Core database. Defaults to the -storage_driver_db value. + -storage_driver_influxdb_bucket + # Use secure connection with database. False by default + -storage_driver_secure + # Writes will be buffered for this duration, and committed to the non memory backends as a single transaction. Default is '60s' + -storage_driver_buffer_duration +``` + +Notes: + +- The bucket (or InfluxDB 3 Core database) must exist before starting cAdvisor. Create it with `influx bucket create` (InfluxDB 2.x) or `influxdb3 create database` (InfluxDB 3 Core). +- The retention policy flag does not apply in token mode; retention is a property of the InfluxDB 2.x bucket. +- InfluxDB 3 Core ignores the organization, but the InfluxDB 2.x API requires the query parameter, so set `-storage_driver_influxdb_org` to any non-empty string. + +### Example: cAdvisor with InfluxDB 2.x + +```console +$ influx bucket create --name cadvisor --org my-org --token $INFLUX_TOKEN +$ docker run \ + --volume=/:/rootfs:ro \ + --volume=/var/run:/var/run:ro \ + --volume=/sys:/sys:ro \ + --volume=/var/lib/docker/:/var/lib/docker:ro \ + --volume=/dev/disk/:/dev/disk:ro \ + --publish=8080:8080 \ + --detach=true \ + --name=cadvisor \ + gcr.io/cadvisor/cadvisor:vlatest \ + -storage_driver=influxdb \ + -storage_driver_host=influxdb:8086 \ + -storage_driver_db=cadvisor \ + -storage_driver_influxdb_auth_token=$INFLUX_TOKEN \ + -storage_driver_influxdb_org=my-org +``` + +### Example: cAdvisor with InfluxDB 3 Core + +```console +$ influxdb3 create database --database cadvisor +$ influxdb3 create token --admin # prints the admin token +$ docker run \ + --volume=/:/rootfs:ro \ + --volume=/var/run:/var/run:ro \ + --volume=/sys:/sys:ro \ + --volume=/var/lib/docker/:/var/lib/docker:ro \ + --volume=/dev/disk/:/dev/disk:ro \ + --publish=8080:8080 \ + --detach=true \ + --name=cadvisor \ + gcr.io/cadvisor/cadvisor:vlatest \ + -storage_driver=influxdb \ + -storage_driver_host=influxdb3:8181 \ + -storage_driver_db=cadvisor \ + -storage_driver_influxdb_auth_token=$INFLUX3_TOKEN \ + -storage_driver_influxdb_org=cadvisor +``` + # Examples [Brian Christner](https://www.brianchristner.io) wrote a detailed post on [setting up Docker monitoring](https://www.brianchristner.io/how-to-setup-docker-monitoring) with cAdvisor and Influxdb. A docker compose configuration for setting up cadvisor-influxdb-grafana can be found [here](https://github.com/dalekurt/docker-monitoring/blob/master/docker-compose.yml).