From 5a526555d7294f7929a29f8b2ff4c1f0e96f9c2c Mon Sep 17 00:00:00 2001 From: Ilya Shevelyov Date: Thu, 23 Jul 2026 12:48:08 +0200 Subject: [PATCH 1/3] feat(autoops_es): add ingest volume enrichment fields (#51404) Extends cat_shards and node_stats metricsets with bulk/dataset/dense-vector snapshot fields and ingest/bulk rate enrichers for resource-utilization insights. Co-Authored-By: Claude Sonnet 4.6 --- .../_meta/test/cat_shards.8.15.3.json | 212 +++++++++++++++--- .../module/autoops_es/cat_shards/cache.go | 32 +++ .../autoops_es/cat_shards/cache_test.go | 125 +++++++++++ .../autoops_es/cat_shards/cat_shards.go | 2 +- .../module/autoops_es/cat_shards/data.go | 5 + .../module/autoops_es/cat_shards/data_test.go | 5 + .../autoops_es/cat_shards/deserialize.go | 4 + .../autoops_es/cat_shards/index_shards.go | 48 ++-- .../cat_shards/index_shards_test.go | 22 ++ .../module/autoops_es/node_stats/cache.go | 4 + .../autoops_es/node_stats/cache_test.go | 56 +++++ .../module/autoops_es/node_stats/data.go | 4 + .../module/autoops_es/node_stats/data_test.go | 20 ++ 13 files changed, 488 insertions(+), 51 deletions(-) diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/_meta/test/cat_shards.8.15.3.json b/x-pack/metricbeat/module/autoops_es/cat_shards/_meta/test/cat_shards.8.15.3.json index 92f2cebdd440..a7c6b8bfbdc6 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/_meta/test/cat_shards.8.15.3.json +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/_meta/test/cat_shards.8.15.3.json @@ -19,7 +19,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -41,7 +45,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -63,7 +71,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -85,7 +97,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "6417", + "dvc": "0" }, { "n": "instance-0000000000", @@ -107,7 +123,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "32784", + "dvc": "0" }, { "n": "instance-0000000000", @@ -129,7 +149,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "24987", + "dvc": "0" }, { "n": "instance-0000000000", @@ -151,7 +175,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -173,7 +201,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -195,7 +227,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -217,7 +253,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -239,7 +279,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -261,7 +305,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -283,7 +331,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -305,7 +357,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -327,7 +383,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -349,7 +409,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -371,7 +435,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -393,7 +461,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -415,7 +487,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -437,7 +513,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -459,7 +539,11 @@ "gmto": "76", "gmti": "10", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "104564", + "dvc": "0" }, { "n": "instance-0000000000", @@ -481,7 +565,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "7012", + "dvc": "0" }, { "n": "instance-0000000000", @@ -503,7 +591,11 @@ "gmto": "1", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "2459155", + "dvc": "0" }, { "n": "instance-0000000000", @@ -525,7 +617,11 @@ "gmto": "14", "gmti": "2", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "308935", + "dvc": "0" }, { "n": "instance-0000000000", @@ -547,7 +643,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "7061", + "dvc": "0" }, { "n": "instance-0000000000", @@ -569,7 +669,11 @@ "gmto": "4", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "31710", + "dvc": "0" }, { "n": "instance-0000000000", @@ -591,7 +695,11 @@ "gmto": "2", "gmti": "4", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "137738", + "dvc": "0" }, { "n": "instance-0000000000", @@ -613,7 +721,11 @@ "gmto": "3", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "360408", + "dvc": "0" }, { "n": "instance-0000000000", @@ -635,7 +747,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "12013", + "dvc": "0" }, { "n": "instance-0000000000", @@ -657,7 +773,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "12516", + "dvc": "0" }, { "n": "instance-0000000000", @@ -679,7 +799,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -701,7 +825,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -723,7 +851,11 @@ "gmto": "0", "gmti": "0", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "227", + "dvc": "0" }, { "n": "instance-0000000000", @@ -745,7 +877,11 @@ "gmto": "31", "gmti": "1", "ur": null, - "ud": null + "ud": null, + "btsi": "1000", + "bto": "10", + "dataset": "98064", + "dvc": "0" }, { "n": "name2", @@ -767,6 +903,10 @@ "gmto": "31", "gmti": "1", "ur": null, - "ud": null + "ud": null, + "btsi": null, + "bto": null, + "dataset": null, + "dvc": null } -] +] \ No newline at end of file diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/cache.go b/x-pack/metricbeat/module/autoops_es/cat_shards/cache.go index ecbf9c98a008..644681e1b3d2 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/cache.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/cache.go @@ -54,6 +54,38 @@ var ( IsUsable: func(obj *NodeIndexShards) bool { return obj.SearchQueryTotal != nil }, WriteValue: func(obj *NodeIndexShards, value float64) { obj.SearchRatePerSecond = &value }, }, + { + CalculateValue: utils.CalculateRate, + ConvertTime: utils.MillisToSeconds, + GetTime: utils.UseTimestamp[*NodeIndexShards], + GetValue: func(obj *NodeIndexShards) int64 { return *obj.DocsCount }, + IsUsable: func(obj *NodeIndexShards) bool { return obj.DocsCount != nil }, + WriteValue: func(obj *NodeIndexShards, value float64) { obj.IngestDocsPerSecond = &value }, + }, + { + CalculateValue: utils.CalculateRate, + ConvertTime: utils.MillisToSeconds, + GetTime: utils.UseTimestamp[*NodeIndexShards], + GetValue: func(obj *NodeIndexShards) int64 { return *obj.SizeInBytes }, + IsUsable: func(obj *NodeIndexShards) bool { return obj.SizeInBytes != nil }, + WriteValue: func(obj *NodeIndexShards, value float64) { obj.IngestBytesPerSecond = &value }, + }, + { + CalculateValue: utils.CalculateRate, + ConvertTime: utils.MillisToSeconds, + GetTime: utils.UseTimestamp[*NodeIndexShards], + GetValue: func(obj *NodeIndexShards) int64 { return *obj.BulkTotalSizeInBytes }, + IsUsable: func(obj *NodeIndexShards) bool { return obj.BulkTotalSizeInBytes != nil }, + WriteValue: func(obj *NodeIndexShards, value float64) { obj.BulkBytesPerSecond = &value }, + }, + { + CalculateValue: utils.CalculateRate, + ConvertTime: utils.MillisToSeconds, + GetTime: utils.UseTimestamp[*NodeIndexShards], + GetValue: func(obj *NodeIndexShards) int64 { return *obj.BulkTotalOperations }, + IsUsable: func(obj *NodeIndexShards) bool { return obj.BulkTotalOperations != nil }, + WriteValue: func(obj *NodeIndexShards, value float64) { obj.BulkOperationsPerSecond = &value }, + }, // LATENCIES: { CalculateValue: utils.CalculateLatency, diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go b/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go index d38cfae84685..bc9045c1719b 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go @@ -58,6 +58,11 @@ func getShard(nodeId string, nodeName string, shardId int32, primary bool, state get_missing_time := []int64{100, 101, 102, 103} get_missing_total := []int64{110, 111, 112, 113} + bulk_total_size_in_bytes := []int64{200, 201, 202, 203} + bulk_total_operations := []int64{300, 301, 302, 303} + total_data_set_size_in_bytes := []int64{10, 11, 12, 13} + dense_vector_count := []int64{0, 0, 0, 0} + shard.docs = &docs[index] shard.store = &store[index] shard.segments_count = &segments_count[index] @@ -70,6 +75,10 @@ func getShard(nodeId string, nodeName string, shardId int32, primary bool, state shard.merges_total_time = &merges_total_time[index] shard.get_missing_time = &get_missing_time[index] shard.get_missing_total = &get_missing_total[index] + shard.bulk_total_size_in_bytes = &bulk_total_size_in_bytes[index] + shard.bulk_total_operations = &bulk_total_operations[index] + shard.total_data_set_size_in_bytes = &total_data_set_size_in_bytes[index] + shard.dense_vector_count = &dense_vector_count[index] return shard } @@ -823,3 +832,119 @@ func TestConvertToNodeIndexShardsWithCache(t *testing.T) { require.Nil(t, myIndexNode3.SearchRatePerSecond) require.Nil(t, myIndexNode3.TimestampDiff) } + +func getNodeIndexShardsWithIngestBulk() map[string]NodeIndexShards { + docsCount := []int64{1, 2} + sizeInBytes := []int64{10, 11} + bulkTotalSizeInBytes := []int64{200, 201} + bulkTotalOperations := []int64{300, 301} + + base := getNodeIndexShards() + + n1 := base["my-index-node_id-node1"] + n1.DocsCount = &docsCount[0] + n1.SizeInBytes = &sizeInBytes[0] + n1.BulkTotalSizeInBytes = &bulkTotalSizeInBytes[0] + n1.BulkTotalOperations = &bulkTotalOperations[0] + base["my-index-node_id-node1"] = n1 + + n2 := base["my-index-node_id-node2"] + n2.DocsCount = &docsCount[1] + n2.SizeInBytes = &sizeInBytes[1] + n2.BulkTotalSizeInBytes = &bulkTotalSizeInBytes[1] + n2.BulkTotalOperations = &bulkTotalOperations[1] + base["my-index-node_id-node2"] = n2 + + return base +} + +func TestEnrichNodeIndexShardsIngestAndBulkRatesWithoutCache(t *testing.T) { + clearCache() + + nodeIndexShardsMap := getNodeIndexShardsWithIngestBulk() + nodeIndexShardsList := enrichNodeIndexShards(nodeIndexShardsMap, map[string]IndexMetadata{}) + + for _, nodeIndexShards := range nodeIndexShardsList { + require.Nil(t, nodeIndexShards.IngestDocsPerSecond) + require.Nil(t, nodeIndexShards.IngestBytesPerSecond) + require.Nil(t, nodeIndexShards.BulkBytesPerSecond) + require.Nil(t, nodeIndexShards.BulkOperationsPerSecond) + } +} + +func TestEnrichNodeIndexShardsIngestAndBulkRatesWithCachedValues(t *testing.T) { + // 10s ago cache + initCache(getNodeIndexShardsWithIngestBulk(), 10) + + nodeIndexShardsMap := getNodeIndexShardsWithIngestBulk() + + for key, nodeIndexShards := range nodeIndexShardsMap { + *nodeIndexShards.DocsCount += 10 // delta=10 → 10/10 = 1/s + *nodeIndexShards.SizeInBytes += 20 // delta=20 → 20/10 = 2/s + *nodeIndexShards.BulkTotalSizeInBytes += 100 // delta=100 → 100/10 = 10/s + *nodeIndexShards.BulkTotalOperations += 50 // delta=50 → 50/10 = 5/s + + nodeIndexShardsMap[key] = nodeIndexShards + } + + nodeIndexShardsList := enrichNodeIndexShards(nodeIndexShardsMap, map[string]IndexMetadata{}) + + for _, nodeIndexShards := range nodeIndexShardsList { + require.EqualValues(t, 1, *nodeIndexShards.IngestDocsPerSecond) + require.EqualValues(t, 2, *nodeIndexShards.IngestBytesPerSecond) + require.EqualValues(t, 10, *nodeIndexShards.BulkBytesPerSecond) + require.EqualValues(t, 5, *nodeIndexShards.BulkOperationsPerSecond) + } +} + +func TestEnrichNodeIndexShardsIngestAndBulkRatesWithNoChange(t *testing.T) { + // 10s ago cache — identical values, expect zero rates + initCache(getNodeIndexShardsWithIngestBulk(), 10) + + nodeIndexShardsMap := getNodeIndexShardsWithIngestBulk() + nodeIndexShardsList := enrichNodeIndexShards(nodeIndexShardsMap, map[string]IndexMetadata{}) + + for _, nodeIndexShards := range nodeIndexShardsList { + require.EqualValues(t, 0, *nodeIndexShards.IngestDocsPerSecond) + require.EqualValues(t, 0, *nodeIndexShards.IngestBytesPerSecond) + require.EqualValues(t, 0, *nodeIndexShards.BulkBytesPerSecond) + require.EqualValues(t, 0, *nodeIndexShards.BulkOperationsPerSecond) + } +} + +func TestEnrichNodeIndexShardsIngestAndBulkRatesWithNewNode(t *testing.T) { + // 10s ago cache with only node1 — node2 is a cache miss + initCache(getNodeIndexShardsWithIngestBulk(), 10) + + nodeIndexShardsMap := getNodeIndexShardsWithIngestBulk() + + for key, nodeIndexShards := range nodeIndexShardsMap { + *nodeIndexShards.DocsCount += 10 + *nodeIndexShards.SizeInBytes += 20 + *nodeIndexShards.BulkTotalSizeInBytes += 100 + *nodeIndexShards.BulkTotalOperations += 50 + nodeIndexShardsMap[key] = nodeIndexShards + } + + newNode := getNodeIndexShardsWithIngestBulk()["my-index-node_id-node2"] + newNode.Index = "my-other-index" + newNode.NodeId = "node3" + newNode.IndexNode = "my-other-index-node_id-node3" + nodeIndexShardsMap["my-other-index-node_id-node3"] = newNode + + nodeIndexShardsList := enrichNodeIndexShards(nodeIndexShardsMap, map[string]IndexMetadata{}) + + for _, nodeIndexShards := range nodeIndexShardsList { + if nodeIndexShards.NodeId != "node3" { + require.EqualValues(t, 1, *nodeIndexShards.IngestDocsPerSecond) + require.EqualValues(t, 2, *nodeIndexShards.IngestBytesPerSecond) + require.EqualValues(t, 10, *nodeIndexShards.BulkBytesPerSecond) + require.EqualValues(t, 5, *nodeIndexShards.BulkOperationsPerSecond) + } else { + require.Nil(t, nodeIndexShards.IngestDocsPerSecond) + require.Nil(t, nodeIndexShards.IngestBytesPerSecond) + require.Nil(t, nodeIndexShards.BulkBytesPerSecond) + require.Nil(t, nodeIndexShards.BulkOperationsPerSecond) + } + } +} diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/cat_shards.go b/x-pack/metricbeat/module/autoops_es/cat_shards/cat_shards.go index 6a2a5a4a0f9a..f0ac10a85958 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/cat_shards.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/cat_shards.go @@ -10,7 +10,7 @@ import ( const ( catShardsMetricSet = "cat_shards" - catShardsPath = "/_cat/shards?s=i&h=n,i,id,s,p,st,d,sto,sc,sqto,sqti,iito,iiti,iif,mt,mtt,gmto,gmti,ur,ud&bytes=b&time=ms&format=json" + catShardsPath = "/_cat/shards?s=i&h=n,i,id,s,p,st,d,sto,sc,sqto,sqti,iito,iiti,iif,mt,mtt,gmto,gmti,ur,ud,btsi,bto,dataset,dvc&bytes=b&time=ms&format=json" resolveIndexPath = "/_resolve/index/*?expand_wildcards=all&filter_path=indices" ) diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/data.go b/x-pack/metricbeat/module/autoops_es/cat_shards/data.go index 1a901f945a9b..fe484ecca4b1 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/data.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/data.go @@ -61,6 +61,11 @@ type JSONShard struct { SearchQueryTime json.Number `json:"sqti"` SearchQueryTotal json.Number `json:"sqto"` + BulkTotalSizeInBytes json.Number `json:"btsi"` + BulkTotalOperations json.Number `json:"bto"` + DatasetSize json.Number `json:"dataset"` + DenseVectorCount json.Number `json:"dvc"` + // only used for unassigned state UnassignedReason *string `json:"ur"` UnassignedDetails *string `json:"ud"` diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/data_test.go b/x-pack/metricbeat/module/autoops_es/cat_shards/data_test.go index 10f856da1252..7c397d7ad274 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/data_test.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/data_test.go @@ -246,6 +246,11 @@ func expectValidParsedDetailedShardsWithCache(t *testing.T, data metricset.Fetch require.Nil(t, myIndexNode2.MergeRatePerSecond) require.Nil(t, myIndexNode2.IndexLatencyInMillis) require.Nil(t, myIndexNode2.MergeLatencyInMillis) + // Replica-only node: primary-only fields are nil, so ingest/bulk rates are nil + require.Nil(t, myIndexNode2.IngestDocsPerSecond) + require.Nil(t, myIndexNode2.IngestBytesPerSecond) + require.Nil(t, myIndexNode2.BulkBytesPerSecond) + require.Nil(t, myIndexNode2.BulkOperationsPerSecond) } // Expect a valid response from Elasticsearch to create N events diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go b/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go index a675e29e5855..7e30edbbbefe 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go @@ -62,6 +62,10 @@ func deserializeShard(jsonShard JSONShard) Shard { shard.segments_count = toInt64(jsonShard.SegmentsCount) shard.search_query_time = toInt64(jsonShard.SearchQueryTime) shard.search_query_total = toInt64(jsonShard.SearchQueryTotal) + shard.bulk_total_size_in_bytes = toInt64(jsonShard.BulkTotalSizeInBytes) + shard.bulk_total_operations = toInt64(jsonShard.BulkTotalOperations) + shard.total_data_set_size_in_bytes = toInt64(jsonShard.DatasetSize) + shard.dense_vector_count = toInt64(jsonShard.DenseVectorCount) return shard } diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go index ecc4f92784fc..8478fc45043b 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go @@ -16,20 +16,24 @@ type Shard struct { state string // optional; set if state != "UNASSIGNED" - docs *int64 - store *int64 - segments_count *int64 - search_query_total *int64 - search_query_time *int64 - indexing_index_total *int64 - indexing_index_time *int64 - indexing_index_failed *int64 - merges_total *int64 - merges_total_time *int64 - get_missing_time *int64 - get_missing_total *int64 - unassigned_reason *string - unassigned_details *string + docs *int64 + store *int64 + segments_count *int64 + search_query_total *int64 + search_query_time *int64 + indexing_index_total *int64 + indexing_index_time *int64 + indexing_index_failed *int64 + merges_total *int64 + merges_total_time *int64 + get_missing_time *int64 + get_missing_total *int64 + bulk_total_size_in_bytes *int64 + bulk_total_operations *int64 + total_data_set_size_in_bytes *int64 + dense_vector_count *int64 + unassigned_reason *string + unassigned_details *string } type AssignedShard struct { @@ -117,6 +121,18 @@ type NodeIndexShards struct { TotalMergesTotal *int64 `json:"total_merges_total"` // includes replicas TotalMergesTotalTime *int64 `json:"total_merges_total_time"` // includes replicas TimestampDiff *int64 `json:"timestamp_diff"` + + // primary-only snapshot fields + BulkTotalSizeInBytes *int64 `json:"bulk_total_size_in_bytes"` + BulkTotalOperations *int64 `json:"bulk_total_operations"` + TotalDataSetSizeInBytes *int64 `json:"total_data_set_size_in_bytes"` + DenseVectorCount *int64 `json:"dense_vector_count"` + + // ingest / bulk rate fields + IngestDocsPerSecond *float64 `json:"ingest_docs_per_second"` + IngestBytesPerSecond *float64 `json:"ingest_bytes_per_second"` + BulkBytesPerSecond *float64 `json:"bulk_bytes_per_second"` + BulkOperationsPerSecond *float64 `json:"bulk_operations_per_second"` } type NodeShardCount struct { @@ -294,6 +310,10 @@ func indexShardsToNodeIndexShards(nodeIndexShardsMap map[string]NodeIndexShards, nodeIndex.IndexingIndexTotalTime = utils.AddInt64OrNull(nodeIndex.IndexingIndexTotalTime, shard.indexing_index_time) nodeIndex.MergesTotal = utils.AddInt64OrNull(nodeIndex.MergesTotal, shard.merges_total) nodeIndex.MergesTotalTime = utils.AddInt64OrNull(nodeIndex.MergesTotalTime, shard.merges_total_time) + nodeIndex.BulkTotalSizeInBytes = utils.AddInt64OrNull(nodeIndex.BulkTotalSizeInBytes, shard.bulk_total_size_in_bytes) + nodeIndex.BulkTotalOperations = utils.AddInt64OrNull(nodeIndex.BulkTotalOperations, shard.bulk_total_operations) + nodeIndex.TotalDataSetSizeInBytes = utils.AddInt64OrNull(nodeIndex.TotalDataSetSizeInBytes, shard.total_data_set_size_in_bytes) + nodeIndex.DenseVectorCount = utils.AddInt64OrNull(nodeIndex.DenseVectorCount, shard.dense_vector_count) } assignedShard := toAssignedShard(shard) diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards_test.go b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards_test.go index e8a43dcf78b3..75ef8635f382 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards_test.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards_test.go @@ -420,6 +420,10 @@ func TestIndexShardsToNodeIndexShardsGreenIndex(t *testing.T) { indexingFailed := []int64{60, 61, 62, 63} indexingTotals := []int64{70, 71, 72, 73} indexingTimes := []int64{74, 75, 76, 77} + bulkSizes := []int64{200, 201, 202, 203} + bulkOps := []int64{300, 301, 302, 303} + datasetSizes := []int64{20, 21, 22, 23} + dvcCounts := []int64{0, 0, 5, 0} shards := []Shard{ {shard: 0, primary: true, node_id: "node1", node_name: "name1"}, @@ -442,6 +446,10 @@ func TestIndexShardsToNodeIndexShardsGreenIndex(t *testing.T) { shard.indexing_index_failed = &indexingFailed[i] shard.indexing_index_total = &indexingTotals[i] shard.indexing_index_time = &indexingTimes[i] + shard.bulk_total_size_in_bytes = &bulkSizes[i] + shard.bulk_total_operations = &bulkOps[i] + shard.total_data_set_size_in_bytes = &datasetSizes[i] + shard.dense_vector_count = &dvcCounts[i] shards[i] = shard } @@ -484,6 +492,10 @@ func TestIndexShardsToNodeIndexShardsGreenIndex(t *testing.T) { require.Equal(t, searchQueryTimes[0], *nodeIndex.SearchQueryTime) require.Equal(t, mergeTotals[0], *nodeIndex.TotalMergesTotal) require.Equal(t, mergeTimes[0], *nodeIndex.TotalMergesTotalTime) + require.Equal(t, bulkSizes[0], *nodeIndex.BulkTotalSizeInBytes) + require.Equal(t, bulkOps[0], *nodeIndex.BulkTotalOperations) + require.Equal(t, datasetSizes[0], *nodeIndex.TotalDataSetSizeInBytes) + require.Equal(t, dvcCounts[0], *nodeIndex.DenseVectorCount) require.Equal(t, 1, len(nodeIndex.AssignShards)) require.Equal(t, 0, len(nodeIndex.InitializingShards)) require.EqualValues(t, 0, nodeIndex.Initializing) @@ -535,6 +547,11 @@ func TestIndexShardsToNodeIndexShardsGreenIndex(t *testing.T) { require.Equal(t, searchQueryTimes[1]+searchQueryTimes[2], *nodeIndex.SearchQueryTime) require.Equal(t, mergeTotals[1]+mergeTotals[2], *nodeIndex.TotalMergesTotal) require.Equal(t, mergeTimes[1]+mergeTimes[2], *nodeIndex.TotalMergesTotalTime) + // primary-only bulk/dataset/dvc fields: only shard[2] (index 2) is primary for node2 + require.Equal(t, bulkSizes[2], *nodeIndex.BulkTotalSizeInBytes) + require.Equal(t, bulkOps[2], *nodeIndex.BulkTotalOperations) + require.Equal(t, datasetSizes[2], *nodeIndex.TotalDataSetSizeInBytes) + require.Equal(t, dvcCounts[2], *nodeIndex.DenseVectorCount) require.Equal(t, 2, len(nodeIndex.AssignShards)) require.Equal(t, 0, len(nodeIndex.InitializingShards)) require.EqualValues(t, 0, nodeIndex.Initializing) @@ -589,6 +606,11 @@ func TestIndexShardsToNodeIndexShardsGreenIndex(t *testing.T) { require.Nil(t, nodeIndex.IndexingIndexTotalTime) require.Nil(t, nodeIndex.MergesTotal) require.Nil(t, nodeIndex.MergesTotalTime) + // replica-only node: no primary-only snapshot fields + require.Nil(t, nodeIndex.BulkTotalSizeInBytes) + require.Nil(t, nodeIndex.BulkTotalOperations) + require.Nil(t, nodeIndex.TotalDataSetSizeInBytes) + require.Nil(t, nodeIndex.DenseVectorCount) require.Equal(t, getMissingTotals[3], *nodeIndex.GetMissingDocTotal) require.Equal(t, getMissingTimes[3], *nodeIndex.GetMissingDocTotalTime) require.Equal(t, searchQueryTotals[3], *nodeIndex.SearchQueryTotal) diff --git a/x-pack/metricbeat/module/autoops_es/node_stats/cache.go b/x-pack/metricbeat/module/autoops_es/node_stats/cache.go index f07a58db03a9..293aa82749bd 100644 --- a/x-pack/metricbeat/module/autoops_es/node_stats/cache.go +++ b/x-pack/metricbeat/module/autoops_es/node_stats/cache.go @@ -17,6 +17,10 @@ var ( createRate("index_rate_per_second", "indices.indexing.index_total"), createRate("merge_rate_per_second", "indices.merges.total"), createRate("search_rate_per_second", "indices.search.query_total"), + createRate("ingest_docs_per_second", "indices.docs.count"), + createRate("ingest_bytes_per_second", "indices.store.size_in_bytes"), + createRate("bulk_bytes_per_second", "indices.bulk.total_size_in_bytes"), + createRate("bulk_operations_per_second", "indices.bulk.total_operations"), // LATENCIES: createLatency("index_latency_in_millis", "indices.indexing.index_total", "indices.indexing.index_time_in_millis"), createLatency("merge_latency_in_millis", "indices.merges.total", "indices.merges.total_time_in_millis"), diff --git a/x-pack/metricbeat/module/autoops_es/node_stats/cache_test.go b/x-pack/metricbeat/module/autoops_es/node_stats/cache_test.go index f85f28518ecc..5ab539be0f24 100644 --- a/x-pack/metricbeat/module/autoops_es/node_stats/cache_test.go +++ b/x-pack/metricbeat/module/autoops_es/node_stats/cache_test.go @@ -31,6 +31,12 @@ func initCache(previousCache map[string]mapstr.M, previousSeconds int64) { func getNodeStatsForNode(nodeIndex int64) mapstr.M { return mapstr.M{ "indices": mapstr.M{ + "docs": mapstr.M{ + "count": 100 + nodeIndex, + }, + "store": mapstr.M{ + "size_in_bytes": 200 + nodeIndex, + }, "indexing": mapstr.M{ "index_failed": 10 + nodeIndex, "index_total": 20 + nodeIndex, @@ -44,6 +50,10 @@ func getNodeStatsForNode(nodeIndex int64) mapstr.M { "query_time_in_millis": 70 + nodeIndex, "query_total": 40 + nodeIndex, }, + "bulk": mapstr.M{ + "total_size_in_bytes": 300 + nodeIndex, + "total_operations": 400 + nodeIndex, + }, }, } } @@ -68,6 +78,10 @@ func TestEnrichNodeStatsWithoutCache(t *testing.T) { require.Nil(t, nodeStatsNode1["index_rate_per_second"]) require.Nil(t, nodeStatsNode1["merge_rate_per_second"]) require.Nil(t, nodeStatsNode1["search_rate_per_second"]) + require.Nil(t, nodeStatsNode1["ingest_docs_per_second"]) + require.Nil(t, nodeStatsNode1["ingest_bytes_per_second"]) + require.Nil(t, nodeStatsNode1["bulk_bytes_per_second"]) + require.Nil(t, nodeStatsNode1["bulk_operations_per_second"]) require.Nil(t, nodeStatsNode1["index_latency_in_millis"]) require.Nil(t, nodeStatsNode1["merge_latency_in_millis"]) require.Nil(t, nodeStatsNode1["search_latency_in_millis"]) @@ -87,6 +101,10 @@ func TestEnrichNodeStatsWithoutCachedValues(t *testing.T) { require.Nil(t, nodeStatsNode1["index_rate_per_second"]) require.Nil(t, nodeStatsNode1["merge_rate_per_second"]) require.Nil(t, nodeStatsNode1["search_rate_per_second"]) + require.Nil(t, nodeStatsNode1["ingest_docs_per_second"]) + require.Nil(t, nodeStatsNode1["ingest_bytes_per_second"]) + require.Nil(t, nodeStatsNode1["bulk_bytes_per_second"]) + require.Nil(t, nodeStatsNode1["bulk_operations_per_second"]) require.Nil(t, nodeStatsNode1["index_latency_in_millis"]) require.Nil(t, nodeStatsNode1["merge_latency_in_millis"]) require.Nil(t, nodeStatsNode1["search_latency_in_millis"]) @@ -106,6 +124,10 @@ func TestEnrichNodeStatsWithCachedValues(t *testing.T) { nodeStats["indices.merges.total_time_in_millis"] = getValue(&nodeStats, "indices.merges.total_time_in_millis") + 40 nodeStats["indices.search.query_total"] = getValue(&nodeStats, "indices.search.query_total") + 60 nodeStats["indices.search.query_time_in_millis"] = getValue(&nodeStats, "indices.search.query_time_in_millis") + 120 + nodeStats["indices.docs.count"] = getValue(&nodeStats, "indices.docs.count") + 10 + nodeStats["indices.store.size_in_bytes"] = getValue(&nodeStats, "indices.store.size_in_bytes") + 20 + nodeStats["indices.bulk.total_size_in_bytes"] = getValue(&nodeStats, "indices.bulk.total_size_in_bytes") + 100 + nodeStats["indices.bulk.total_operations"] = getValue(&nodeStats, "indices.bulk.total_operations") + 50 nodeStatsMap[key] = nodeStats } @@ -124,6 +146,10 @@ func TestEnrichNodeStatsWithCachedValues(t *testing.T) { require.EqualValues(t, 3, nodeStats["index_failed_rate_per_second"]) require.EqualValues(t, 4, nodeStats["merge_rate_per_second"]) require.EqualValues(t, 6, nodeStats["search_rate_per_second"]) + require.EqualValues(t, 1, nodeStats["ingest_docs_per_second"]) + require.EqualValues(t, 2, nodeStats["ingest_bytes_per_second"]) + require.EqualValues(t, 10, nodeStats["bulk_bytes_per_second"]) + require.EqualValues(t, 5, nodeStats["bulk_operations_per_second"]) // latencies require.EqualValues(t, 0.5, nodeStats["index_latency_in_millis"]) require.EqualValues(t, 1, nodeStats["merge_latency_in_millis"]) @@ -169,6 +195,12 @@ func TestEnrichNodeStatsSearchLatencyClampsToInterval(t *testing.T) { "index latency below interval should not be clamped") require.InDelta(t, 10, nodeStats["merge_latency_in_millis"], 0.01, "merge latency below interval should not be clamped") + + // No increments for ingest/bulk counters → rates are zero + require.EqualValues(t, 0, nodeStats["ingest_docs_per_second"]) + require.EqualValues(t, 0, nodeStats["ingest_bytes_per_second"]) + require.EqualValues(t, 0, nodeStats["bulk_bytes_per_second"]) + require.EqualValues(t, 0, nodeStats["bulk_operations_per_second"]) } } @@ -192,6 +224,10 @@ func TestEnrichNodeStatsWithCachedValuesWithNoChange(t *testing.T) { require.EqualValues(t, 0, nodeStats["index_failed_rate_per_second"]) require.EqualValues(t, 0, nodeStats["merge_rate_per_second"]) require.EqualValues(t, 0, nodeStats["search_rate_per_second"]) + require.EqualValues(t, 0, nodeStats["ingest_docs_per_second"]) + require.EqualValues(t, 0, nodeStats["ingest_bytes_per_second"]) + require.EqualValues(t, 0, nodeStats["bulk_bytes_per_second"]) + require.EqualValues(t, 0, nodeStats["bulk_operations_per_second"]) // latencies require.EqualValues(t, 0, nodeStats["index_latency_in_millis"]) require.EqualValues(t, 0, nodeStats["merge_latency_in_millis"]) @@ -221,6 +257,10 @@ func TestEnrichNodeStatsWithCachedValuesWithHoles(t *testing.T) { nodeStats["indices.merges.total_time_in_millis"] = getValue(&nodeStats, "indices.merges.total_time_in_millis") + 40 nodeStats["indices.search.query_total"] = getValue(&nodeStats, "indices.search.query_total") + 60 nodeStats["indices.search.query_time_in_millis"] = getValue(&nodeStats, "indices.search.query_time_in_millis") + 120 + nodeStats["indices.docs.count"] = getValue(&nodeStats, "indices.docs.count") + 10 + nodeStats["indices.store.size_in_bytes"] = getValue(&nodeStats, "indices.store.size_in_bytes") + 20 + nodeStats["indices.bulk.total_size_in_bytes"] = getValue(&nodeStats, "indices.bulk.total_size_in_bytes") + 100 + nodeStats["indices.bulk.total_operations"] = getValue(&nodeStats, "indices.bulk.total_operations") + 50 nodeStatsMap[key] = nodeStats } @@ -244,6 +284,10 @@ func TestEnrichNodeStatsWithCachedValuesWithHoles(t *testing.T) { require.EqualValues(t, 3, nodeStats["index_failed_rate_per_second"]) require.EqualValues(t, 4, nodeStats["merge_rate_per_second"]) require.EqualValues(t, 6, nodeStats["search_rate_per_second"]) + require.EqualValues(t, 1, nodeStats["ingest_docs_per_second"]) + require.EqualValues(t, 2, nodeStats["ingest_bytes_per_second"]) + require.EqualValues(t, 10, nodeStats["bulk_bytes_per_second"]) + require.EqualValues(t, 5, nodeStats["bulk_operations_per_second"]) // latencies if key == "node2" { @@ -270,6 +314,10 @@ func TestEnrichNodeIndexShardsWithCachedValuesWithNewNodeAndIndex(t *testing.T) nodeStats["indices.merges.total_time_in_millis"] = getValue(&nodeStats, "indices.merges.total_time_in_millis") + 40 nodeStats["indices.search.query_total"] = getValue(&nodeStats, "indices.search.query_total") + 60 nodeStats["indices.search.query_time_in_millis"] = getValue(&nodeStats, "indices.search.query_time_in_millis") + 120 + nodeStats["indices.docs.count"] = getValue(&nodeStats, "indices.docs.count") + 10 + nodeStats["indices.store.size_in_bytes"] = getValue(&nodeStats, "indices.store.size_in_bytes") + 20 + nodeStats["indices.bulk.total_size_in_bytes"] = getValue(&nodeStats, "indices.bulk.total_size_in_bytes") + 100 + nodeStats["indices.bulk.total_operations"] = getValue(&nodeStats, "indices.bulk.total_operations") + 50 nodeStatsMap[key] = nodeStats } @@ -295,6 +343,10 @@ func TestEnrichNodeIndexShardsWithCachedValuesWithNewNodeAndIndex(t *testing.T) require.EqualValues(t, 3, nodeStats["index_failed_rate_per_second"]) require.EqualValues(t, 4, nodeStats["merge_rate_per_second"]) require.EqualValues(t, 6, nodeStats["search_rate_per_second"]) + require.EqualValues(t, 1, nodeStats["ingest_docs_per_second"]) + require.EqualValues(t, 2, nodeStats["ingest_bytes_per_second"]) + require.EqualValues(t, 10, nodeStats["bulk_bytes_per_second"]) + require.EqualValues(t, 5, nodeStats["bulk_operations_per_second"]) // latencies require.EqualValues(t, 0.5, nodeStats["index_latency_in_millis"]) require.EqualValues(t, 1, nodeStats["merge_latency_in_millis"]) @@ -305,6 +357,10 @@ func TestEnrichNodeIndexShardsWithCachedValuesWithNewNodeAndIndex(t *testing.T) require.Nil(t, nodeStats["index_rate_per_second"]) require.Nil(t, nodeStats["merge_rate_per_second"]) require.Nil(t, nodeStats["search_rate_per_second"]) + require.Nil(t, nodeStats["ingest_docs_per_second"]) + require.Nil(t, nodeStats["ingest_bytes_per_second"]) + require.Nil(t, nodeStats["bulk_bytes_per_second"]) + require.Nil(t, nodeStats["bulk_operations_per_second"]) require.Nil(t, nodeStats["index_latency_in_millis"]) require.Nil(t, nodeStats["merge_latency_in_millis"]) require.Nil(t, nodeStats["search_latency_in_millis"]) diff --git a/x-pack/metricbeat/module/autoops_es/node_stats/data.go b/x-pack/metricbeat/module/autoops_es/node_stats/data.go index 0c1b989e3cbc..6a5d3fe9c018 100644 --- a/x-pack/metricbeat/module/autoops_es/node_stats/data.go +++ b/x-pack/metricbeat/module/autoops_es/node_stats/data.go @@ -99,6 +99,10 @@ var ( "total_size_bytes": c.Int("total_size_bytes", s.IgnoreAllErrors), }, c.DictOptional), }, c.DictOptional), + "bulk": c.Dict("bulk", s.Schema{ + "total_size_in_bytes": c.Int("total_size_in_bytes", s.IgnoreAllErrors), + "total_operations": c.Int("total_operations", s.IgnoreAllErrors), + }, c.DictOptional), }, c.DictOptional), "os": c.Dict("os", s.Schema{ "cpu": c.Dict("cpu", s.Schema{ diff --git a/x-pack/metricbeat/module/autoops_es/node_stats/data_test.go b/x-pack/metricbeat/module/autoops_es/node_stats/data_test.go index c340c4276bcb..318913a00c21 100644 --- a/x-pack/metricbeat/module/autoops_es/node_stats/data_test.go +++ b/x-pack/metricbeat/module/autoops_es/node_stats/data_test.go @@ -107,6 +107,8 @@ func expectValidParsedDetailed(t *testing.T, data metricset.FetcherData[NodesSta require.EqualValues(t, 175109606, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.search.query_total")) require.EqualValues(t, 3464297906, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.search.query_time_in_millis")) require.EqualValues(t, 5358, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.segments.count")) + require.EqualValues(t, 277737431206550, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.bulk.total_size_in_bytes")) + require.EqualValues(t, 17452535488, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.bulk.total_operations")) require.EqualValues(t, 32, auto_ops_testing.GetObjectValue(node1MetricSet, "thread_pool.write.threads")) require.EqualValues(t, 24175874622, auto_ops_testing.GetObjectValue(node1MetricSet, "thread_pool.write.completed")) require.EqualValues(t, 1, auto_ops_testing.GetObjectValue(node1MetricSet, "thread_pool.snapshot.threads")) @@ -120,6 +122,8 @@ func expectValidParsedDetailed(t *testing.T, data metricset.FetcherData[NodesSta require.Equal(t, false, node1MetricSet["is_elected_master"]) require.ElementsMatch(t, []string{"data_content", "data_hot", "ingest", "master", "remote_cluster_client", "transform"}, node1MetricSet["roles"]) require.EqualValues(t, 2902603, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.docs.count")) + require.EqualValues(t, 593809157, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.bulk.total_size_in_bytes")) + require.EqualValues(t, 646725, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.bulk.total_operations")) require.EqualValues(t, 2818360, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.dense_vector.count")) require.EqualValues(t, 4510292000, auto_ops_testing.GetObjectValue(node1MetricSet, "indices.dense_vector.off_heap.total_size_bytes")) } @@ -139,6 +143,10 @@ func expectValidParsedDetailedWithNoCache(t *testing.T, data metricset.FetcherDa require.Nil(t, node1MetricSet["index_rate_per_second"]) require.Nil(t, node1MetricSet["merge_rate_per_second"]) require.Nil(t, node1MetricSet["search_rate_per_second"]) + require.Nil(t, node1MetricSet["ingest_docs_per_second"]) + require.Nil(t, node1MetricSet["ingest_bytes_per_second"]) + require.Nil(t, node1MetricSet["bulk_bytes_per_second"]) + require.Nil(t, node1MetricSet["bulk_operations_per_second"]) require.Nil(t, node1MetricSet["index_latency_in_millis"]) require.Nil(t, node1MetricSet["merge_latency_in_millis"]) require.Nil(t, node1MetricSet["search_latency_in_millis"]) @@ -159,6 +167,18 @@ func expectValidParsedDetailedWithCache(t *testing.T, data metricset.FetcherData require.NotNil(t, node1MetricSet["index_latency_in_millis"]) require.NotNil(t, node1MetricSet["merge_latency_in_millis"]) require.NotNil(t, node1MetricSet["search_latency_in_millis"]) + + // ingest rates use docs/store which are present on all ES versions + require.NotNil(t, node1MetricSet["ingest_docs_per_second"]) + require.NotNil(t, node1MetricSet["ingest_bytes_per_second"]) + // bulk rates require indices.bulk which is only present on ES 8+ + if data.Version == "7.17.0" { + require.Nil(t, node1MetricSet["bulk_bytes_per_second"]) + require.Nil(t, node1MetricSet["bulk_operations_per_second"]) + } else { + require.NotNil(t, node1MetricSet["bulk_bytes_per_second"]) + require.NotNil(t, node1MetricSet["bulk_operations_per_second"]) + } } // Expect a valid response from Elasticsearch to create N events From d82ae6e130bb51fd295dc873ceee497ae7c235af Mon Sep 17 00:00:00 2001 From: Ilya Shevelyov Date: Thu, 23 Jul 2026 15:56:35 +0200 Subject: [PATCH 2/3] rename JSONShard.DatasetSize to TotalDataSetSizeInBytes for consistency Co-Authored-By: Claude Sonnet 4.6 --- .../autoops_es/cat_shards/cache_test.go | 4 +-- .../module/autoops_es/cat_shards/data.go | 8 ++--- .../autoops_es/cat_shards/deserialize.go | 2 +- .../autoops_es/cat_shards/index_shards.go | 34 +++++++++---------- .../module/autoops_es/node_stats/data.go | 2 +- 5 files changed, 25 insertions(+), 25 deletions(-) diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go b/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go index bc9045c1719b..f17f48ef3dc0 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/cache_test.go @@ -879,8 +879,8 @@ func TestEnrichNodeIndexShardsIngestAndBulkRatesWithCachedValues(t *testing.T) { nodeIndexShardsMap := getNodeIndexShardsWithIngestBulk() for key, nodeIndexShards := range nodeIndexShardsMap { - *nodeIndexShards.DocsCount += 10 // delta=10 → 10/10 = 1/s - *nodeIndexShards.SizeInBytes += 20 // delta=20 → 20/10 = 2/s + *nodeIndexShards.DocsCount += 10 // delta=10 → 10/10 = 1/s + *nodeIndexShards.SizeInBytes += 20 // delta=20 → 20/10 = 2/s *nodeIndexShards.BulkTotalSizeInBytes += 100 // delta=100 → 100/10 = 10/s *nodeIndexShards.BulkTotalOperations += 50 // delta=50 → 50/10 = 5/s diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/data.go b/x-pack/metricbeat/module/autoops_es/cat_shards/data.go index fe484ecca4b1..74154bebc6ed 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/data.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/data.go @@ -61,10 +61,10 @@ type JSONShard struct { SearchQueryTime json.Number `json:"sqti"` SearchQueryTotal json.Number `json:"sqto"` - BulkTotalSizeInBytes json.Number `json:"btsi"` - BulkTotalOperations json.Number `json:"bto"` - DatasetSize json.Number `json:"dataset"` - DenseVectorCount json.Number `json:"dvc"` + BulkTotalSizeInBytes json.Number `json:"btsi"` + BulkTotalOperations json.Number `json:"bto"` + TotalDataSetSizeInBytes json.Number `json:"dataset"` + DenseVectorCount json.Number `json:"dvc"` // only used for unassigned state UnassignedReason *string `json:"ur"` diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go b/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go index 7e30edbbbefe..1aac2c5ca5ba 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/deserialize.go @@ -64,7 +64,7 @@ func deserializeShard(jsonShard JSONShard) Shard { shard.search_query_total = toInt64(jsonShard.SearchQueryTotal) shard.bulk_total_size_in_bytes = toInt64(jsonShard.BulkTotalSizeInBytes) shard.bulk_total_operations = toInt64(jsonShard.BulkTotalOperations) - shard.total_data_set_size_in_bytes = toInt64(jsonShard.DatasetSize) + shard.total_data_set_size_in_bytes = toInt64(jsonShard.TotalDataSetSizeInBytes) shard.dense_vector_count = toInt64(jsonShard.DenseVectorCount) return shard diff --git a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go index 8478fc45043b..3c7bd4dcff13 100644 --- a/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go +++ b/x-pack/metricbeat/module/autoops_es/cat_shards/index_shards.go @@ -16,24 +16,24 @@ type Shard struct { state string // optional; set if state != "UNASSIGNED" - docs *int64 - store *int64 - segments_count *int64 - search_query_total *int64 - search_query_time *int64 - indexing_index_total *int64 - indexing_index_time *int64 - indexing_index_failed *int64 - merges_total *int64 - merges_total_time *int64 - get_missing_time *int64 - get_missing_total *int64 - bulk_total_size_in_bytes *int64 - bulk_total_operations *int64 + docs *int64 + store *int64 + segments_count *int64 + search_query_total *int64 + search_query_time *int64 + indexing_index_total *int64 + indexing_index_time *int64 + indexing_index_failed *int64 + merges_total *int64 + merges_total_time *int64 + get_missing_time *int64 + get_missing_total *int64 + bulk_total_size_in_bytes *int64 + bulk_total_operations *int64 total_data_set_size_in_bytes *int64 - dense_vector_count *int64 - unassigned_reason *string - unassigned_details *string + dense_vector_count *int64 + unassigned_reason *string + unassigned_details *string } type AssignedShard struct { diff --git a/x-pack/metricbeat/module/autoops_es/node_stats/data.go b/x-pack/metricbeat/module/autoops_es/node_stats/data.go index 6a5d3fe9c018..60c0cfb7eceb 100644 --- a/x-pack/metricbeat/module/autoops_es/node_stats/data.go +++ b/x-pack/metricbeat/module/autoops_es/node_stats/data.go @@ -227,7 +227,7 @@ type ClusterStateMasterNode struct { } type NodesStats struct { - Nodes map[string]map[string]interface{} `json:"nodes"` + Nodes map[string]map[string]any `json:"nodes"` } // Get the elected master node's ID From cdda877802dada0a594794b2469e6406421c4cad Mon Sep 17 00:00:00 2001 From: Ilya Shevelyov Date: Thu, 23 Jul 2026 17:14:12 +0200 Subject: [PATCH 3/3] add changelog fragment for ingest volume enrichment fields Co-Authored-By: Claude Sonnet 4.6 --- .../1784819509-autoops-es-ingest-volume-enriched-fields.yaml | 3 +++ 1 file changed, 3 insertions(+) create mode 100644 changelog/fragments/1784819509-autoops-es-ingest-volume-enriched-fields.yaml diff --git a/changelog/fragments/1784819509-autoops-es-ingest-volume-enriched-fields.yaml b/changelog/fragments/1784819509-autoops-es-ingest-volume-enriched-fields.yaml new file mode 100644 index 000000000000..9ab8d68d7678 --- /dev/null +++ b/changelog/fragments/1784819509-autoops-es-ingest-volume-enriched-fields.yaml @@ -0,0 +1,3 @@ +kind: enhancement +summary: Add ingest volume enrichment fields to autoops_es module +component: metricbeat