master
go 208 lines 4.16 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package hdfs
4
5 import (
6 "encoding/json"
7 "errors"
8 "fmt"
9 "strings"
10
11 "github.com/netdata/netdata/go/plugins/pkg/stm"
12 "github.com/netdata/netdata/go/plugins/pkg/web"
13 )
14
15 const (
16 dataNodeType = "DataNode"
17 nameNodeType = "NameNode"
18 )
19
20 func (c *Collector) collect() (map[string]int64, error) {
21 req, err := web.NewHTTPRequest(c.RequestConfig)
22 if err != nil {
23 return nil, fmt.Errorf("failed to create HTTP request: %v", err)
24 }
25
26 var raw rawJMX
27 if err := web.DoHTTP(c.httpClient).RequestJSON(req, &raw); err != nil {
28 return nil, err
29 }
30
31 if raw.isEmpty() {
32 return nil, errors.New("empty response")
33 }
34
35 mx := c.collectRawJMX(raw)
36
37 return stm.ToMap(mx), nil
38 }
39
40 func (c *Collector) determineNodeType() (string, error) {
41 req, err := web.NewHTTPRequest(c.RequestConfig)
42 if err != nil {
43 return "", fmt.Errorf("failed to create HTTP request: %v", err)
44 }
45
46 var raw rawJMX
47 if err := web.DoHTTP(c.httpClient).RequestJSON(req, &raw); err != nil {
48 return "", err
49 }
50
51 if raw.isEmpty() {
52 return "", errors.New("empty response")
53 }
54
55 jvm := raw.findJvm()
56 if jvm == nil {
57 return "", errors.New("couldn't find jvm in response")
58 }
59
60 v, ok := jvm["tag.ProcessName"]
61 if !ok {
62 return "", errors.New("couldn't find process name in JvmMetrics")
63 }
64
65 t := strings.Trim(string(v), "\"")
66 if t == nameNodeType || t == dataNodeType {
67 return t, nil
68 }
69 return "", errors.New("unknown node type")
70 }
71
72 func (c *Collector) collectRawJMX(raw rawJMX) *metrics {
73 var mx metrics
74 switch c.nodeType {
75 default:
76 panic(fmt.Sprintf("unsupported node type : '%s'", c.nodeType))
77 case nameNodeType:
78 c.collectNameNode(&mx, raw)
79 case dataNodeType:
80 c.collectDataNode(&mx, raw)
81 }
82 return &mx
83 }
84
85 func (c *Collector) collectNameNode(mx *metrics, raw rawJMX) {
86 if err := c.collectJVM(mx, raw); err != nil {
87 c.Debugf("error on collecting jvm : %v", err)
88 }
89
90 if err := c.collectRPCActivity(mx, raw); err != nil {
91 c.Debugf("error on collecting rpc activity : %v", err)
92 }
93
94 if err := c.collectFSNameSystem(mx, raw); err != nil {
95 c.Debugf("error on collecting fs name system : %v", err)
96 }
97 }
98
99 func (c *Collector) collectDataNode(mx *metrics, raw rawJMX) {
100 if err := c.collectJVM(mx, raw); err != nil {
101 c.Debugf("error on collecting jvm : %v", err)
102 }
103
104 if err := c.collectRPCActivity(mx, raw); err != nil {
105 c.Debugf("error on collecting rpc activity : %v", err)
106 }
107
108 if err := c.collectFSDatasetState(mx, raw); err != nil {
109 c.Debugf("error on collecting fs dataset state : %v", err)
110 }
111
112 if err := c.collectDataNodeActivity(mx, raw); err != nil {
113 c.Debugf("error on collecting datanode activity state : %v", err)
114 }
115 }
116
117 func (c *Collector) collectJVM(mx *metrics, raw rawJMX) error {
118 v := raw.findJvm()
119 if v == nil {
120 return nil
121 }
122
123 var jvm jvmMetrics
124 err := writeJSONTo(&jvm, v)
125 if err != nil {
126 return err
127 }
128
129 mx.Jvm = &jvm
130 return nil
131 }
132
133 func (c *Collector) collectRPCActivity(mx *metrics, raw rawJMX) error {
134 v := raw.findRPCActivity()
135 if v == nil {
136 return nil
137 }
138
139 var rpc rpcActivityMetrics
140 err := writeJSONTo(&rpc, v)
141 if err != nil {
142 return err
143 }
144
145 mx.Rpc = &rpc
146 return nil
147 }
148
149 func (c *Collector) collectFSNameSystem(mx *metrics, raw rawJMX) error {
150 v := raw.findFSNameSystem()
151 if v == nil {
152 return nil
153 }
154
155 var fs fsNameSystemMetrics
156 err := writeJSONTo(&fs, v)
157 if err != nil {
158 return err
159 }
160
161 fs.CapacityUsed = fs.CapacityDfsUsed + fs.CapacityUsedNonDFS
162
163 mx.FSNameSystem = &fs
164 return nil
165 }
166
167 func (c *Collector) collectFSDatasetState(mx *metrics, raw rawJMX) error {
168 v := raw.findFSDatasetState()
169 if v == nil {
170 return nil
171 }
172
173 var fs fsDatasetStateMetrics
174 err := writeJSONTo(&fs, v)
175 if err != nil {
176 return err
177 }
178
179 fs.CapacityUsed = fs.Capacity - fs.Remaining
180 fs.CapacityUsedNonDFS = fs.CapacityUsed - fs.DfsUsed
181
182 mx.FSDatasetState = &fs
183 return nil
184 }
185
186 func (c *Collector) collectDataNodeActivity(mx *metrics, raw rawJMX) error {
187 v := raw.findDataNodeActivity()
188 if v == nil {
189 return nil
190 }
191
192 var dna dataNodeActivityMetrics
193 err := writeJSONTo(&dna, v)
194 if err != nil {
195 return err
196 }
197
198 mx.DataNodeActivity = &dna
199 return nil
200 }
201
202 func writeJSONTo(dst, src any) error {
203 b, err := json.Marshal(src)
204 if err != nil {
205 return err
206 }
207 return json.Unmarshal(b, dst)
208 }