Monitoring Cloudflare's edge network with Prometheus
drive.google.com
drive.google.com
Personally I'm very happy that the open source world is adopting something derived from Borgmon rather than something derived from its supposed "replacement".
It wasn't meant to be commentary on Prometheus(which I quite like) at all :)
Borgmon may be dead in the eyes of some people, but I know for a fact that it's still the only thing monitoring core and critical systems.
Most of the problem with Borgmon, IMO, is the cruft that has built up over the decade+, and neglect due to the Google pattern of "The new thing that doesn't work, and the old thing that is deprecated.".
The difficulty at Google is that developers are rewarded for writing new and shiny from scratch, rather than fix the old but working systems.
This isn't always a problem, as some good things can come out of starting from scratch. But sometimes they throw out too many of the good ideas, in an attempt to be fancy and new.
For Google-scale orgs or infrastructure needs. Most everyone else in the world does not need Google scale tools.
Most of them do run their own datacenters, sometimes in numerous locations, they have massive and extremely complex IT systems in place.
I won't weight in on that debate. But you can think of Prometheus as an experiment to decide the issue: it is very similar to Borgmon, but has a cleaner language.
It's config language is less crazy (Python based) and operates globally.
https://www.youtube.com/watch?v=LlvJdK1xsl4
Edit: Monarch config isn't sane, it's just different and at least not in the crazy languages that borgmon uses.
As for Monarch, it's a very different beast. For one, it stores all its rules in a protocol buffer format, so it's more structured. But then you have to write Python code that generates the protocol buffers and pushes them to storage. It looks similar but not the same as the ad-hoc query language. I wouldn't go as far as calling it sane.
It is also a service and it's optimized for Google's network architecture with datacenter local and global nodes and the language itself is aware of this distinction and some computations are done locally, others globally and so on.
For your local monitoring needs (or even global ones, if you're willing to put in the effort), Prometheus is a solid choice.
The equivalent easy to setup and use, with ALL the features working out of the box, SaaS standard is https://www.datadoghq.com/ or potentially Google Stack driver if you are on Google Cloud.
[1] https://news.ycombinator.com/item?id=15315028
Disclaimer: No relation to either org.
-- Winston Churchill if he worked at Google.
Such as Prometheus :) Even some teams in Google use it.
1. As already mentioned, the macro system (and the fact that its use is basically required to set up basic monitoring) has quite a steep learning curve. Prometheus doesn't have that and I personally would prefer it did.
2. It's not a service, so you have to set up your own instance, configure it, maintain it. In many engineers' mind this is just another hurdle in front of them launching their service.
3. As a software engineer (particularly new to Google) you might not expect to have to do ops work and carry a pager yourself.
Out of all three, only (1) is a valid reason to hate on borgmon. That and the language itself, which is almost a 1:1 match with Prometheus, are very different from your regular programming language. But given the choice between flat, simple metrics (which is what most monitoring systems give you) and the ability to have arbitrary dimensions and be able to work with them to build useful alerts and dashboards and troubleshoot quickly, I (again personal opinion) will always go with the latter.
What makes its pull-based mechanism superior to push-based ones like statsd?
And using exporters sounds clunky - instead of directly querying a metric and sending it to your metrics collector, you have an intermediate component which exposes them for collection.
With push, you need to know where to send the data. So you depend on a fixed configuration. With pull, you can essentially do a periodic nmap sweep to discover and update the monitoring data sources.
That's one less thing to go wrong.
Let's say I have a service that's only used once per day, at a random point in the day. I need to know if the service is available all the time, even if it's not used. With pull, I'll get back an 'OK' e.g. every minute. You could say that we should implement heartbeats, but again, then this service needs to be configurable to pick up changes in the monitoring infrastructure, then I need to have a separate check for the presence of heartbeats and the values of the samples..
What if the service spins up on demand when it's needed, does its job, and then exits. You want to track how often it starts, how long it runs, how long the job was in the queue, how much CPU/RAM/IO was consumed, etc.
With push that's all part of the cleanup code in the process. With pull it seems like you could completely miss that the process was even ever run?
The push vs pull debate is relevant for longer running daemons, not batch jobs.
I found that people who're new to larger-scale monitoring favor push, because that's somehow more intuitive; but pull really works very well, it's not clunky at all.
The metrics collection needs to be done by a local agent installed on the system, for the reason you gave, that's the only place the data is available.
The metrics storage is somewhere else.
Prometheus does the storage. For the collection, you still have to install collectd, statsd or similar on your hosts. Sure, prometheus could do a HTTP check remotely, but that won't get cover much of anything.
The Prometheus pattern with central things like a database (e.g. monitoring Postgres) is to let the exporter (the thing that acts as an intermediary to expose metrics for Prometheus to fetch) run anywhere it wants. It absolutely does not need to be a "local agent".
(In fact, if you're using something like Kubernetes, the only thing that needs local access to a node is the exporter that exposes node-specific metrics. Everything else can chat over the network.)
The benefit here is that if you have 10 Postgres databases, you can still run just a single exporter and have it extract data from all of them. Or you can run one exporter per database; there's conceptually little difference.
On Kubernetes, we usually run an exporter as a sidecar container, which means it can talk to Postgres or whatever on localhost and just live alongside the process that it's exporting metrics for, and we rely on Prometheus' automatic discovery to make Prometheus pull from it. Start a new Postgres instance and its metrics are almost immediately fed into Prometheus.
You don't need collectd etc. with Prometheus. There are exporters around for just about anything.
They can be called agents, exporters, collectors, whatever, the name is not important, the design pattern is.
A system that would be exclusively pull-based from the prometheus server does not work practically.
Unless you take the presence of node_exporter to invalidate the premise (which would be stupid) a system that's exclusively pull-based is entirely practical.
Pulling has a few technical benefits, though. For one, only the puller needs to know what's being monitored; the thing being monitored can therefore be exceedingly simple, dumb and passive. Statsd is similarly simple in that it's just local UDP broadcast, of course, which leads to the next point:
Another benefit is that it allows better fine-grained control over when metrics gathering is done, and what. Since Prometheus best practices dictate that metrics should be computed at pull time, it means you can fine-tune collection intervals to specific metrics, and this can be done centrally. And since you only pull from what you have, it means there can't be a rogue agent somewhere that's spewing out data (i.e. what a sibling comment calls "authorative sources").
But to understand why pull is a better model, you have to understand Google's/Prometheus's "observer/reactor" mindset towards large-scale computing; it's just easier to scale up with this model. Consider an application that implements some kind of REST API. You want metrics for things like the total number of requests served, which you'll sample now and then. You add an endpoint /metrics running on port 9100. Then you tell Prometheus to scrape (pull from) http://example.com:9100/metrics. So far so good.
The beauty of the model arises when you involve a dynamic orchestration like Kubernetes. Now we're running the app on Kubernetes, which means the app can run on many nodes, across many clusters, at the same time; it will have a lot of different IPs (one IP per instance) that are completely dynamic. Instead of adding a rule to scrape a specific URL, you tell Prometheus to ask Kubernetes for all services and then use that information to figure out the endpoint. This dynamic discovery means that as you take apps up and down, Prometheus will automatically update its list of endpoints and scrape them. Equally importing, Prometheus goes to the source of the data at any given time. The services are already scaled up; there's no corresponding metrics collection to scale up, other than in the internal machinery of Prometheus' scraping system.
In other words, Prometheus is observing the cluster and reacting to changes in it to reconfigure it self. This isn't exactly new, but it's core to Google's/Prometheus's way of thinking about applications and services, which has subseqently coloured the whole Kubernetes culture. Instead of configuring the chess pieces, you let the board inspect the chess pieces and configure itself. You want the individual, lower-level apps to be as mundane as possible, let the behavioural signals flow upstream, and let the higher-level pieces make decisions.
This dovetails nicely with the observational data model you need for monitoring, anyway: First you collect the data, then you check the data, then you report anomalies within the data. For example, if you're measuring some number that can go critically high, you don't make the application issue a warning if it goes above a threshold; rather, you collect the data from the application as a raw number, then perform calculations (e.g. max over the last N mins, sum over the last N mins, total count, etc.) that you compare against the threshold.
In practice, implementing a metrics endpoint is exceedingly simple, and you get used to "just writing another exporter". I've written a lot of exporters, and while this initially struck me as heavyweight and clunky, my mindset is now that an HTTP listener is actually more lightweight than an "imperative" pusher script.
Obviously endpoints already need to know how to contact all sorts of services they depend on. So it's not like you're "saving" anything by not telling them "PrometheusIP = X".
Let's say you want to cleanly shut-down some instances of your endpoint. They are holding connection stats & request counts that you don't want to lose. With push the endpoint can close its connection handler, finish any outstanding requests, push final stats, and then exit. With pull are you supposed to just sit and wait until a Pull happens before the process can exit?
* Many installations run multiple Prometheus servers for redundancy, so to start, it'd have to be multiple IPs.
* They would also need auth credentials.
* They'd need retry/failure logic with backoff to prevent dogpiling.
* Clients would have to be careful to resolve the name, not cache the DNS lookup, in order to always resolve Prometheus to the right IP.
* If Prometheus moves, every pusher has to be updated.
* Since Prometheus wouldn't know about pushers, it wouldn't know if a push has failed. As Prometheus is pull-based, you can detect actual failure, not just absence of data.
There's a lot to be said for Prometheus' principle of baking exporters into individual, completely self-encapsulated programs — as opposed to things like collectd, diamond, Munin, Nagios etc. that collect a lot of stuff into a single, possibly plugin-based, system.
Don't forget, a lot of exporters come with third-party software. You want those programs to have as little config as possible. If I release an open-source app (let's say, a search engine), I can include a /metrics handler, and users who deploy my app can just point their Prometheus at it. It's enticingly simple.
As for graceful shutdown: The default pull frequency is 15 seconds, and you can increase it if you want to avoid losing metrics. Prometheus is designed not to deal with extremely fine-grained metrics; losing a few requests due to a shutdown shouldn't matter in the big picture. But for metrics that are sensitive, it's easy enough to bake them into some stateful store anyway (Redis or etcd, for example), or computing them in real time from stateful data (e.g. SQL). For example, if you have some kind of e-commerce order system, it's better if the exporter produces the numbers by issuing a query against the transaction tables, rather than maintaining RAM counters of dollars and cents.
But yes it is possible to miss some requests if a node goes down without Prometheus collecting the latest stats.
But as the parent said, if you need such totals it might be better to store them persistently. Also I do not know a scenario where the total number of requests will trigger an alert.
The rate() function allows for this, you'll get the right answer on average.
>"if you're measuring some number that can go critically high, you don't make the application issue a warning if it goes above a threshold; rather, you collect the data from the application as a raw number, then perform calculations (e.g. max over the last N mins, sum over the last N mins, total count, etc.) that you compare against the threshold."
With standard statsd, you're sending a message per event, this might be fine for one or two servers, with a trickle of traffic. But when you've got 100 servers handling 100s to 1000s of requests per second, we're talking about 10-100k/sec. This is a large amount of load to be directly sending to your network.
With statsd, you could buffer, but now you're losing the benefit of the live event stream.
With Prometheus, we just say keep it all in memory, localy, and as thread-local as possible. This way the cost to update an event is very, very tiny. Much smaller than even a log line. This allows every debug level log line to be recorded as an event counter. For example, the Prometheus Go library only requires ~15 nanoseconds of CPU time to update a single counter.
This cheapness allows for sprinkling metrics everywhere in your code with little worry that it'll be a performance problem.
On the topic of exporters, stand-alone exporters are only necessary where the existing code doesn't directly support Prometheus's metrics format. The good news here is that we're working with several other metrics systems in order to create a common format for polling and pushing metrics. Yes, insert XKCD joke here. :-)
Once we have a better common metrics format, think SNMP but not crazy, "exporters" will be unnecessary.
Prometheus is not intentionally designed as a long term cold storage option for metrics. You _can_ store metrics for as long as your storage allows, but Prometheus is not going to replicate or manage that data to prevent long term degradation. Depending on your long term needs, the preferred pattern is to roll data off Prometheus into a metrics store that better handles data over the scope of months. Rolling data off of Prometheus is done with an exporter and documented [here](https://prometheus.io/docs/operating/integrations/#remote-en...).
In most use operational use cases I've come across, we've only needed about a month of data, so keeping it in Prometheus was kosher. YMMV.
Right now I use graphite but would love something that also handles replication/redundancy with a good query language & enough performance to also use it to fetch the underlying data used to render front-end graphs for users.
The big questions to ask are: - what data is of interest - what are the requested ingest patterns (how is data getting into this system? How frequently? What rate of ingest is expected?) - what are the requested query patterns (who is doing queries? What do those queries look like specifically? Are people querying over unbounded time ranges or are people doing more focused queries? Do queries regularly involve aggregations or not?) - what is the requested SLA (i.e. how much partial down time is okay? How much full downtime is okay? What kind of query response time do you need to target?) - what resources are/will be available (money, man power, compute, storage, tech on hand)
It's possible that a time series data store might not be the correct system choice once these questions are answered. It's not unusual to see data split or copied into multiple systems to answer all requirements.
So I would say, the system is 5 minute ingest of about 1 million metrics. This is spread out over a half dozen locations, each which currently records in their own silo.
And aggregate metric is calculated with a 5 minute lag, which reads all the just-written data points and aggregates into sum-totals which are themselves stored and cached in one place. This is another million metrics basically stored separate from the rest.
But it doesn't really change the character of the system. In the end I'm trying to; Write batches of mostly numeric data Queries over time against those numbers Aggregate data over different blocks of time; 5 minute, hourly, daily, monthly Store it efficiently Ensure redundancy, integrity
Seems like a simple and common enough problem to have been reasonably "solved" for orders of magnitude higher scale than I'm operating at.
At reasonable scale, I've seen people get really far with pure graphite setups by utilizing tools like [carbon-c-relay](https://github.com/grobian/carbon-c-relay), [carbonate](https://github.com/graphite-project/carbonate), and high integrity filesystems like zfs underneath. It's a very hands on operation though. Things like growing the cluster are hard to do without downtime.
If constant growth and uptime is a concern, something like openTSDB might be a great choice. The complexity of setting up a Hadoop + HBase cluster is a pretty big upfront cost, but man is this thing the cockroach of time series data stores. Adding storage is just growing the HBase cluster. Querying across years of data is pretty simple and straightforward. For the complexity involved, openTSDB is worth it.
Prometheus is not recommended for uses where 100% accuracy is required, see https://prometheus.io/docs/introduction/overview/#when-does-...?
I'd also be wary of using Graphite for such a use case. For billing a more traditional database is probably best.
Is "planet scale" better than "web scale"?
You may not be aware but yes, cloudflare is a significant internet company.
I am aware that Cloudflare is a CDN, most CDNs are substantial internet companies. However I've worked at a couple of CDNs and they all throw a lot of marketing numbers around.
https://i2.wp.com/stratusly.com/wp-content/uploads/2017/02/t...
And sounding cringe-worthy in the process :)