2016-01-31 21:45:44 +00:00
|
|
|
package mesos
|
|
|
|
|
|
|
|
import (
|
|
|
|
"encoding/json"
|
|
|
|
"errors"
|
|
|
|
"io/ioutil"
|
2016-02-03 02:31:39 +00:00
|
|
|
"log"
|
2016-01-31 21:45:44 +00:00
|
|
|
"net"
|
|
|
|
"net/http"
|
2018-02-08 02:36:38 +00:00
|
|
|
"net/url"
|
2016-02-09 22:49:30 +00:00
|
|
|
"strconv"
|
2018-02-08 02:36:38 +00:00
|
|
|
"strings"
|
2016-02-02 01:17:38 +00:00
|
|
|
"sync"
|
2016-02-29 16:52:58 +00:00
|
|
|
"time"
|
2016-01-31 21:45:44 +00:00
|
|
|
|
|
|
|
"github.com/influxdata/telegraf"
|
2018-05-04 23:33:23 +00:00
|
|
|
"github.com/influxdata/telegraf/internal/tls"
|
2016-01-31 21:45:44 +00:00
|
|
|
"github.com/influxdata/telegraf/plugins/inputs"
|
2016-02-09 22:57:48 +00:00
|
|
|
jsonparser "github.com/influxdata/telegraf/plugins/parsers/json"
|
2016-01-31 21:45:44 +00:00
|
|
|
)
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
type Role string
|
|
|
|
|
|
|
|
const (
|
|
|
|
MASTER Role = "master"
|
|
|
|
SLAVE = "slave"
|
|
|
|
)
|
|
|
|
|
2016-01-31 21:45:44 +00:00
|
|
|
type Mesos struct {
|
2016-02-09 22:49:30 +00:00
|
|
|
Timeout int
|
2016-02-11 00:06:51 +00:00
|
|
|
Masters []string
|
2016-02-11 00:54:05 +00:00
|
|
|
MasterCols []string `toml:"master_collections"`
|
2016-06-09 10:33:14 +00:00
|
|
|
Slaves []string
|
|
|
|
SlaveCols []string `toml:"slave_collections"`
|
2016-09-23 19:39:59 +00:00
|
|
|
//SlaveTasks bool
|
2018-05-04 23:33:23 +00:00
|
|
|
tls.ClientConfig
|
2018-02-08 02:36:38 +00:00
|
|
|
|
|
|
|
initialized bool
|
|
|
|
client *http.Client
|
|
|
|
masterURLs []*url.URL
|
|
|
|
slaveURLs []*url.URL
|
2016-01-31 21:45:44 +00:00
|
|
|
}
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
var allMetrics = map[Role][]string{
|
|
|
|
MASTER: []string{"resources", "master", "system", "agents", "frameworks", "tasks", "messages", "evqueue", "registrar"},
|
|
|
|
SLAVE: []string{"resources", "agent", "system", "executors", "tasks", "messages"},
|
2016-02-04 01:46:20 +00:00
|
|
|
}
|
|
|
|
|
2016-02-09 23:05:58 +00:00
|
|
|
var sampleConfig = `
|
2016-06-09 10:33:14 +00:00
|
|
|
## Timeout, in ms.
|
2016-02-09 23:05:58 +00:00
|
|
|
timeout = 100
|
2016-06-09 10:33:14 +00:00
|
|
|
## A list of Mesos masters.
|
2018-02-08 02:36:38 +00:00
|
|
|
masters = ["http://localhost:5050"]
|
2016-06-09 10:33:14 +00:00
|
|
|
## Master metrics groups to be collected, by default, all enabled.
|
2016-03-31 23:50:24 +00:00
|
|
|
master_collections = [
|
|
|
|
"resources",
|
|
|
|
"master",
|
|
|
|
"system",
|
2016-06-09 10:33:14 +00:00
|
|
|
"agents",
|
2016-03-31 23:50:24 +00:00
|
|
|
"frameworks",
|
2016-06-09 10:33:14 +00:00
|
|
|
"tasks",
|
2016-03-31 23:50:24 +00:00
|
|
|
"messages",
|
|
|
|
"evqueue",
|
|
|
|
"registrar",
|
|
|
|
]
|
2016-06-09 10:33:14 +00:00
|
|
|
## A list of Mesos slaves, default is []
|
|
|
|
# slaves = []
|
|
|
|
## Slave metrics groups to be collected, by default, all enabled.
|
|
|
|
# slave_collections = [
|
|
|
|
# "resources",
|
|
|
|
# "agent",
|
|
|
|
# "system",
|
|
|
|
# "executors",
|
|
|
|
# "tasks",
|
|
|
|
# "messages",
|
|
|
|
# ]
|
2018-02-08 02:36:38 +00:00
|
|
|
|
2018-05-04 23:33:23 +00:00
|
|
|
## Optional TLS Config
|
|
|
|
# tls_ca = "/etc/telegraf/ca.pem"
|
|
|
|
# tls_cert = "/etc/telegraf/cert.pem"
|
|
|
|
# tls_key = "/etc/telegraf/key.pem"
|
|
|
|
## Use TLS but skip chain & host verification
|
2018-02-08 02:36:38 +00:00
|
|
|
# insecure_skip_verify = false
|
2016-02-09 23:05:58 +00:00
|
|
|
`
|
|
|
|
|
2016-02-03 02:31:39 +00:00
|
|
|
// SampleConfig returns a sample configuration block
|
|
|
|
func (m *Mesos) SampleConfig() string {
|
|
|
|
return sampleConfig
|
|
|
|
}
|
|
|
|
|
|
|
|
// Description just returns a short description of the Mesos plugin
|
|
|
|
func (m *Mesos) Description() string {
|
|
|
|
return "Telegraf plugin for gathering metrics from N Mesos masters"
|
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
func parseURL(s string, role Role) (*url.URL, error) {
|
|
|
|
if !strings.HasPrefix(s, "http://") && !strings.HasPrefix(s, "https://") {
|
|
|
|
host, port, err := net.SplitHostPort(s)
|
|
|
|
// no port specified
|
|
|
|
if err != nil {
|
|
|
|
host = s
|
|
|
|
switch role {
|
|
|
|
case MASTER:
|
|
|
|
port = "5050"
|
|
|
|
case SLAVE:
|
|
|
|
port = "5051"
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
s = "http://" + host + ":" + port
|
|
|
|
log.Printf("W! [inputs.mesos] Using %q as connection URL; please update your configuration to use an URL", s)
|
|
|
|
}
|
|
|
|
|
|
|
|
return url.Parse(s)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *Mesos) initialize() error {
|
2016-06-09 10:33:14 +00:00
|
|
|
if len(m.MasterCols) == 0 {
|
|
|
|
m.MasterCols = allMetrics[MASTER]
|
|
|
|
}
|
|
|
|
|
|
|
|
if len(m.SlaveCols) == 0 {
|
|
|
|
m.SlaveCols = allMetrics[SLAVE]
|
|
|
|
}
|
|
|
|
|
|
|
|
if m.Timeout == 0 {
|
2018-02-08 02:36:38 +00:00
|
|
|
log.Println("I! [inputs.mesos] Missing timeout value, setting default value (100ms)")
|
2016-06-09 10:33:14 +00:00
|
|
|
m.Timeout = 100
|
|
|
|
}
|
2018-02-08 02:36:38 +00:00
|
|
|
|
|
|
|
rawQuery := "timeout=" + strconv.Itoa(m.Timeout) + "ms"
|
|
|
|
|
|
|
|
m.masterURLs = make([]*url.URL, 0, len(m.Masters))
|
|
|
|
for _, master := range m.Masters {
|
|
|
|
u, err := parseURL(master, MASTER)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
u.RawQuery = rawQuery
|
|
|
|
m.masterURLs = append(m.masterURLs, u)
|
|
|
|
}
|
|
|
|
|
|
|
|
m.slaveURLs = make([]*url.URL, 0, len(m.Slaves))
|
|
|
|
for _, slave := range m.Slaves {
|
|
|
|
u, err := parseURL(slave, SLAVE)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
u.RawQuery = rawQuery
|
|
|
|
m.slaveURLs = append(m.slaveURLs, u)
|
|
|
|
}
|
|
|
|
|
|
|
|
client, err := m.createHttpClient()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
m.client = client
|
|
|
|
|
|
|
|
return nil
|
2016-06-09 10:33:14 +00:00
|
|
|
}
|
|
|
|
|
2016-02-09 23:05:58 +00:00
|
|
|
// Gather() metrics from given list of Mesos Masters
|
2016-02-03 02:31:39 +00:00
|
|
|
func (m *Mesos) Gather(acc telegraf.Accumulator) error {
|
2018-02-08 02:36:38 +00:00
|
|
|
if !m.initialized {
|
|
|
|
err := m.initialize()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
m.initialized = true
|
|
|
|
}
|
2016-02-03 02:31:39 +00:00
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
var wg sync.WaitGroup
|
2016-02-03 02:31:39 +00:00
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
for _, master := range m.masterURLs {
|
2016-02-03 02:31:39 +00:00
|
|
|
wg.Add(1)
|
2018-02-08 02:36:38 +00:00
|
|
|
go func(master *url.URL) {
|
|
|
|
acc.AddError(m.gatherMainMetrics(master, MASTER, acc))
|
2016-06-09 10:33:14 +00:00
|
|
|
wg.Done()
|
|
|
|
return
|
2018-02-08 02:36:38 +00:00
|
|
|
}(master)
|
2016-06-09 10:33:14 +00:00
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
for _, slave := range m.slaveURLs {
|
2016-06-09 10:33:14 +00:00
|
|
|
wg.Add(1)
|
2018-02-08 02:36:38 +00:00
|
|
|
go func(slave *url.URL) {
|
|
|
|
acc.AddError(m.gatherMainMetrics(slave, SLAVE, acc))
|
2016-06-09 10:33:14 +00:00
|
|
|
wg.Done()
|
|
|
|
return
|
2018-02-08 02:36:38 +00:00
|
|
|
}(slave)
|
2016-06-09 10:33:14 +00:00
|
|
|
|
2016-09-23 19:39:59 +00:00
|
|
|
// if !m.SlaveTasks {
|
|
|
|
// continue
|
|
|
|
// }
|
2016-06-09 10:33:14 +00:00
|
|
|
|
2016-09-23 19:39:59 +00:00
|
|
|
// wg.Add(1)
|
|
|
|
// go func(c string) {
|
2018-02-08 02:36:38 +00:00
|
|
|
// acc.AddError(m.gatherSlaveTaskMetrics(slave, acc))
|
2016-09-23 19:39:59 +00:00
|
|
|
// wg.Done()
|
|
|
|
// return
|
|
|
|
// }(v)
|
2016-02-03 02:31:39 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
func (m *Mesos) createHttpClient() (*http.Client, error) {
|
2018-05-04 23:33:23 +00:00
|
|
|
tlsCfg, err := m.ClientConfig.TLSConfig()
|
2018-02-08 02:36:38 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
client := &http.Client{
|
|
|
|
Transport: &http.Transport{
|
|
|
|
Proxy: http.ProxyFromEnvironment,
|
|
|
|
TLSClientConfig: tlsCfg,
|
|
|
|
},
|
|
|
|
Timeout: 4 * time.Second,
|
|
|
|
}
|
|
|
|
|
|
|
|
return client, nil
|
|
|
|
}
|
|
|
|
|
2016-02-09 23:05:58 +00:00
|
|
|
// metricsDiff() returns set names for removal
|
2016-06-09 10:33:14 +00:00
|
|
|
func metricsDiff(role Role, w []string) []string {
|
2016-02-04 01:46:20 +00:00
|
|
|
b := []string{}
|
|
|
|
s := make(map[string]bool)
|
|
|
|
|
|
|
|
if len(w) == 0 {
|
|
|
|
return b
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, v := range w {
|
|
|
|
s[v] = true
|
|
|
|
}
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
for _, d := range allMetrics[role] {
|
2016-02-04 01:46:20 +00:00
|
|
|
if _, ok := s[d]; !ok {
|
|
|
|
b = append(b, d)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return b
|
|
|
|
}
|
|
|
|
|
2016-02-09 23:05:58 +00:00
|
|
|
// masterBlocks serves as kind of metrics registry groupping them in sets
|
2016-06-09 10:33:14 +00:00
|
|
|
func getMetrics(role Role, group string) []string {
|
2016-01-31 21:45:44 +00:00
|
|
|
var m map[string][]string
|
|
|
|
|
|
|
|
m = make(map[string][]string)
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
if role == MASTER {
|
|
|
|
m["resources"] = []string{
|
|
|
|
"master/cpus_percent",
|
|
|
|
"master/cpus_used",
|
|
|
|
"master/cpus_total",
|
|
|
|
"master/cpus_revocable_percent",
|
|
|
|
"master/cpus_revocable_total",
|
|
|
|
"master/cpus_revocable_used",
|
|
|
|
"master/disk_percent",
|
|
|
|
"master/disk_used",
|
|
|
|
"master/disk_total",
|
|
|
|
"master/disk_revocable_percent",
|
|
|
|
"master/disk_revocable_total",
|
|
|
|
"master/disk_revocable_used",
|
|
|
|
"master/gpus_percent",
|
|
|
|
"master/gpus_used",
|
|
|
|
"master/gpus_total",
|
|
|
|
"master/gpus_revocable_percent",
|
|
|
|
"master/gpus_revocable_total",
|
|
|
|
"master/gpus_revocable_used",
|
|
|
|
"master/mem_percent",
|
|
|
|
"master/mem_used",
|
|
|
|
"master/mem_total",
|
|
|
|
"master/mem_revocable_percent",
|
|
|
|
"master/mem_revocable_total",
|
|
|
|
"master/mem_revocable_used",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["master"] = []string{
|
|
|
|
"master/elected",
|
|
|
|
"master/uptime_secs",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["system"] = []string{
|
|
|
|
"system/cpus_total",
|
|
|
|
"system/load_15min",
|
|
|
|
"system/load_5min",
|
|
|
|
"system/load_1min",
|
|
|
|
"system/mem_free_bytes",
|
|
|
|
"system/mem_total_bytes",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["agents"] = []string{
|
|
|
|
"master/slave_registrations",
|
|
|
|
"master/slave_removals",
|
|
|
|
"master/slave_reregistrations",
|
|
|
|
"master/slave_shutdowns_scheduled",
|
|
|
|
"master/slave_shutdowns_canceled",
|
|
|
|
"master/slave_shutdowns_completed",
|
|
|
|
"master/slaves_active",
|
|
|
|
"master/slaves_connected",
|
|
|
|
"master/slaves_disconnected",
|
|
|
|
"master/slaves_inactive",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["frameworks"] = []string{
|
|
|
|
"master/frameworks_active",
|
|
|
|
"master/frameworks_connected",
|
|
|
|
"master/frameworks_disconnected",
|
|
|
|
"master/frameworks_inactive",
|
|
|
|
"master/outstanding_offers",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["tasks"] = []string{
|
|
|
|
"master/tasks_error",
|
|
|
|
"master/tasks_failed",
|
|
|
|
"master/tasks_finished",
|
|
|
|
"master/tasks_killed",
|
|
|
|
"master/tasks_lost",
|
|
|
|
"master/tasks_running",
|
|
|
|
"master/tasks_staging",
|
|
|
|
"master/tasks_starting",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["messages"] = []string{
|
|
|
|
"master/invalid_executor_to_framework_messages",
|
|
|
|
"master/invalid_framework_to_executor_messages",
|
|
|
|
"master/invalid_status_update_acknowledgements",
|
|
|
|
"master/invalid_status_updates",
|
|
|
|
"master/dropped_messages",
|
|
|
|
"master/messages_authenticate",
|
|
|
|
"master/messages_deactivate_framework",
|
|
|
|
"master/messages_decline_offers",
|
|
|
|
"master/messages_executor_to_framework",
|
|
|
|
"master/messages_exited_executor",
|
|
|
|
"master/messages_framework_to_executor",
|
|
|
|
"master/messages_kill_task",
|
|
|
|
"master/messages_launch_tasks",
|
|
|
|
"master/messages_reconcile_tasks",
|
|
|
|
"master/messages_register_framework",
|
|
|
|
"master/messages_register_slave",
|
|
|
|
"master/messages_reregister_framework",
|
|
|
|
"master/messages_reregister_slave",
|
|
|
|
"master/messages_resource_request",
|
|
|
|
"master/messages_revive_offers",
|
|
|
|
"master/messages_status_update",
|
|
|
|
"master/messages_status_update_acknowledgement",
|
|
|
|
"master/messages_unregister_framework",
|
|
|
|
"master/messages_unregister_slave",
|
|
|
|
"master/messages_update_slave",
|
|
|
|
"master/recovery_slave_removals",
|
|
|
|
"master/slave_removals/reason_registered",
|
|
|
|
"master/slave_removals/reason_unhealthy",
|
|
|
|
"master/slave_removals/reason_unregistered",
|
|
|
|
"master/valid_framework_to_executor_messages",
|
|
|
|
"master/valid_status_update_acknowledgements",
|
|
|
|
"master/valid_status_updates",
|
|
|
|
"master/task_lost/source_master/reason_invalid_offers",
|
|
|
|
"master/task_lost/source_master/reason_slave_removed",
|
|
|
|
"master/task_lost/source_slave/reason_executor_terminated",
|
|
|
|
"master/valid_executor_to_framework_messages",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["evqueue"] = []string{
|
|
|
|
"master/event_queue_dispatches",
|
|
|
|
"master/event_queue_http_requests",
|
|
|
|
"master/event_queue_messages",
|
|
|
|
}
|
|
|
|
|
|
|
|
m["registrar"] = []string{
|
|
|
|
"registrar/state_fetch_ms",
|
|
|
|
"registrar/state_store_ms",
|
|
|
|
"registrar/state_store_ms/max",
|
|
|
|
"registrar/state_store_ms/min",
|
|
|
|
"registrar/state_store_ms/p50",
|
|
|
|
"registrar/state_store_ms/p90",
|
|
|
|
"registrar/state_store_ms/p95",
|
|
|
|
"registrar/state_store_ms/p99",
|
|
|
|
"registrar/state_store_ms/p999",
|
|
|
|
"registrar/state_store_ms/p9999",
|
|
|
|
}
|
|
|
|
} else if role == SLAVE {
|
|
|
|
m["resources"] = []string{
|
|
|
|
"slave/cpus_percent",
|
|
|
|
"slave/cpus_used",
|
|
|
|
"slave/cpus_total",
|
|
|
|
"slave/cpus_revocable_percent",
|
|
|
|
"slave/cpus_revocable_total",
|
|
|
|
"slave/cpus_revocable_used",
|
|
|
|
"slave/disk_percent",
|
|
|
|
"slave/disk_used",
|
|
|
|
"slave/disk_total",
|
|
|
|
"slave/disk_revocable_percent",
|
|
|
|
"slave/disk_revocable_total",
|
|
|
|
"slave/disk_revocable_used",
|
|
|
|
"slave/gpus_percent",
|
|
|
|
"slave/gpus_used",
|
|
|
|
"slave/gpus_total",
|
|
|
|
"slave/gpus_revocable_percent",
|
|
|
|
"slave/gpus_revocable_total",
|
|
|
|
"slave/gpus_revocable_used",
|
|
|
|
"slave/mem_percent",
|
|
|
|
"slave/mem_used",
|
|
|
|
"slave/mem_total",
|
|
|
|
"slave/mem_revocable_percent",
|
|
|
|
"slave/mem_revocable_total",
|
|
|
|
"slave/mem_revocable_used",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m["agent"] = []string{
|
|
|
|
"slave/registered",
|
|
|
|
"slave/uptime_secs",
|
|
|
|
}
|
|
|
|
|
|
|
|
m["system"] = []string{
|
|
|
|
"system/cpus_total",
|
|
|
|
"system/load_15min",
|
|
|
|
"system/load_5min",
|
|
|
|
"system/load_1min",
|
|
|
|
"system/mem_free_bytes",
|
|
|
|
"system/mem_total_bytes",
|
|
|
|
}
|
|
|
|
|
|
|
|
m["executors"] = []string{
|
|
|
|
"containerizer/mesos/container_destroy_errors",
|
|
|
|
"slave/container_launch_errors",
|
|
|
|
"slave/executors_preempted",
|
|
|
|
"slave/frameworks_active",
|
|
|
|
"slave/executor_directory_max_allowed_age_secs",
|
|
|
|
"slave/executors_registering",
|
|
|
|
"slave/executors_running",
|
|
|
|
"slave/executors_terminated",
|
|
|
|
"slave/executors_terminating",
|
|
|
|
"slave/recovery_errors",
|
|
|
|
}
|
|
|
|
|
|
|
|
m["tasks"] = []string{
|
|
|
|
"slave/tasks_failed",
|
|
|
|
"slave/tasks_finished",
|
|
|
|
"slave/tasks_killed",
|
|
|
|
"slave/tasks_lost",
|
|
|
|
"slave/tasks_running",
|
|
|
|
"slave/tasks_staging",
|
|
|
|
"slave/tasks_starting",
|
|
|
|
}
|
|
|
|
|
|
|
|
m["messages"] = []string{
|
|
|
|
"slave/invalid_framework_messages",
|
|
|
|
"slave/invalid_status_updates",
|
|
|
|
"slave/valid_framework_messages",
|
|
|
|
"slave/valid_status_updates",
|
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
}
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
ret, ok := m[group]
|
2016-01-31 21:45:44 +00:00
|
|
|
|
|
|
|
if !ok {
|
2017-11-01 00:00:06 +00:00
|
|
|
log.Printf("I! [mesos] Unknown %s metrics group: %s\n", role, group)
|
2016-02-03 02:31:39 +00:00
|
|
|
return []string{}
|
2016-01-31 21:45:44 +00:00
|
|
|
}
|
|
|
|
|
2016-02-03 02:31:39 +00:00
|
|
|
return ret
|
2016-01-31 21:45:44 +00:00
|
|
|
}
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
func (m *Mesos) filterMetrics(role Role, metrics *map[string]interface{}) {
|
2016-02-03 02:31:39 +00:00
|
|
|
var ok bool
|
2016-06-09 10:33:14 +00:00
|
|
|
var selectedMetrics []string
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
if role == MASTER {
|
|
|
|
selectedMetrics = m.MasterCols
|
|
|
|
} else if role == SLAVE {
|
|
|
|
selectedMetrics = m.SlaveCols
|
|
|
|
}
|
2016-02-02 01:17:38 +00:00
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
for _, k := range metricsDiff(role, selectedMetrics) {
|
|
|
|
for _, v := range getMetrics(role, k) {
|
|
|
|
if _, ok = (*metrics)[v]; ok {
|
|
|
|
delete((*metrics), v)
|
2016-02-04 01:46:20 +00:00
|
|
|
}
|
2016-02-03 02:31:39 +00:00
|
|
|
}
|
2016-02-02 01:17:38 +00:00
|
|
|
}
|
2016-01-31 21:45:44 +00:00
|
|
|
}
|
|
|
|
|
2016-08-30 14:25:29 +00:00
|
|
|
// TaskStats struct for JSON API output /monitor/statistics
|
|
|
|
type TaskStats struct {
|
|
|
|
ExecutorID string `json:"executor_id"`
|
|
|
|
FrameworkID string `json:"framework_id"`
|
|
|
|
Statistics map[string]interface{} `json:"statistics"`
|
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
func (m *Mesos) gatherSlaveTaskMetrics(u *url.URL, acc telegraf.Accumulator) error {
|
2016-08-30 14:25:29 +00:00
|
|
|
var metrics []TaskStats
|
2016-06-09 10:33:14 +00:00
|
|
|
|
|
|
|
tags := map[string]string{
|
2018-02-08 02:36:38 +00:00
|
|
|
"server": u.Hostname(),
|
|
|
|
"url": urlTag(u),
|
2016-06-09 10:33:14 +00:00
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
resp, err := m.client.Get(withPath(u, "/monitor/statistics").String())
|
2016-06-09 10:33:14 +00:00
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
data, err := ioutil.ReadAll(resp.Body)
|
|
|
|
resp.Body.Close()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if err = json.Unmarshal([]byte(data), &metrics); err != nil {
|
|
|
|
return errors.New("Error decoding JSON response")
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, task := range metrics {
|
2016-08-30 14:25:29 +00:00
|
|
|
tags["framework_id"] = task.FrameworkID
|
2016-06-09 10:33:14 +00:00
|
|
|
|
|
|
|
jf := jsonparser.JSONFlattener{}
|
2016-08-30 14:25:29 +00:00
|
|
|
err = jf.FlattenJSON("", task.Statistics)
|
2016-06-09 10:33:14 +00:00
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2016-08-30 15:44:12 +00:00
|
|
|
|
2016-08-30 14:25:29 +00:00
|
|
|
timestamp := time.Unix(int64(jf.Fields["timestamp"].(float64)), 0)
|
2016-08-30 15:44:12 +00:00
|
|
|
jf.Fields["executor_id"] = task.ExecutorID
|
2016-06-09 10:33:14 +00:00
|
|
|
|
2016-08-30 14:25:29 +00:00
|
|
|
acc.AddFields("mesos_tasks", jf.Fields, tags, timestamp)
|
2016-06-09 10:33:14 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
func withPath(u *url.URL, path string) *url.URL {
|
|
|
|
c := *u
|
|
|
|
c.Path = path
|
|
|
|
return &c
|
|
|
|
}
|
|
|
|
|
|
|
|
func urlTag(u *url.URL) string {
|
|
|
|
c := *u
|
|
|
|
c.Path = ""
|
|
|
|
c.User = nil
|
|
|
|
c.RawQuery = ""
|
|
|
|
return c.String()
|
|
|
|
}
|
|
|
|
|
2016-02-03 02:31:39 +00:00
|
|
|
// This should not belong to the object
|
2018-02-08 02:36:38 +00:00
|
|
|
func (m *Mesos) gatherMainMetrics(u *url.URL, role Role, acc telegraf.Accumulator) error {
|
2016-01-31 21:45:44 +00:00
|
|
|
var jsonOut map[string]interface{}
|
|
|
|
|
|
|
|
tags := map[string]string{
|
2018-02-08 02:36:38 +00:00
|
|
|
"server": u.Hostname(),
|
|
|
|
"url": urlTag(u),
|
2016-06-09 10:33:14 +00:00
|
|
|
"role": string(role),
|
2016-02-09 22:49:30 +00:00
|
|
|
}
|
|
|
|
|
2018-02-08 02:36:38 +00:00
|
|
|
resp, err := m.client.Get(withPath(u, "/metrics/snapshot").String())
|
2016-01-31 21:45:44 +00:00
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
data, err := ioutil.ReadAll(resp.Body)
|
|
|
|
resp.Body.Close()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if err = json.Unmarshal([]byte(data), &jsonOut); err != nil {
|
|
|
|
return errors.New("Error decoding JSON response")
|
|
|
|
}
|
|
|
|
|
2016-06-09 10:33:14 +00:00
|
|
|
m.filterMetrics(role, &jsonOut)
|
2016-01-31 21:45:44 +00:00
|
|
|
|
2016-02-09 22:57:48 +00:00
|
|
|
jf := jsonparser.JSONFlattener{}
|
2016-01-31 21:45:44 +00:00
|
|
|
|
|
|
|
err = jf.FlattenJSON("", jsonOut)
|
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2016-08-30 14:25:29 +00:00
|
|
|
if role == MASTER {
|
|
|
|
if jf.Fields["master/elected"] != 0.0 {
|
|
|
|
tags["state"] = "leader"
|
|
|
|
} else {
|
|
|
|
tags["state"] = "standby"
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2016-01-31 21:45:44 +00:00
|
|
|
acc.AddFields("mesos", jf.Fields, tags)
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func init() {
|
|
|
|
inputs.Add("mesos", func() telegraf.Input {
|
|
|
|
return &Mesos{}
|
|
|
|
})
|
|
|
|
}
|