Building Netflix's Distributed Tracing Infrastructure
netflixtechblog.com
netflixtechblog.com
In my case though, I think people are going about mostly wrongly -- my ideal observability tool is:
- structured logs as the source of truth. Metrics and tracing are actually just special cases of any normal log stream. For applications it doesn't get much simpler than log-to-stdout/stderr, and the log shipping tools we have today can pick out/filter and bucket the metrics/traces/regular logs that come through (or do it on a node-local collector). You don't even have to necessarily trace context very well if you know that the logs from roughly 1-5 seconds are likely relevant and you can filter for signal (warnings/errors/etc) there.
- all-in-one but pluggable with reasonable defaults, for easy administration (because you don't really care how team X does their observability, you just want your own little instance with your own data in it)
- federated (rather than worrying about scaling one mega large service, why not run little ones that can move queries to data and come back?)
- Histograms, because distribution/bucketing is the fastest way to solve a lot of these problems
- schema driven (ex. sprinkle in some JSONLD)
- standardized RED/USE for everything (so your load balancer, network card, CPU, as low/high as you dare/care to go vertically)
To bring all these "features" together, I think the UX you'd put on top of that is incident driven and focused -- context is what matters for solving non-obvious problems, and getting the right context as quickly as possible usually has to do with distance-to-incident. In addition to this, the real next-level goal is to never have an un-analyzed/tagged incident, and soon you can start recognizing (via simple rules, at least) problems and suggesting solutions, or identifying problematic commits before they land (in your infra-as-code repo, or your code-as-code repos).
Theoretically perhaps, but for high-volume metrics (think counters in a tight loop), you really want some form or pre-aggregation in memory before flushing out to disk or network.
If you look at the Prometheus model for example, there's no flushing to disk or network at all, but rather a Prometheus server scraping from hosts on demand.
If you develop an open source project and one of your goals is "I sure wish my users could use Datadog instead of being forced to use an open source thing!" then OpenTelemetry is a project you should keep an eye on. The goal is to make the instrumentation agnostic to the underlying provider.
> By 2017, open source projects like Open-Tracing and Open-Zipkin were mature enough for use in polyglot runtime environments at Netflix.
I don't think OpenCensus even existed in 2017, looks to have been announced in 2018:
InfluxDB is more recent and quite limited. Things like sharding is not supported and they stated it would never be supported except in a paid edition when they make one.
Prometheus is more recent. Similar story with scaling. They changed storage formats and rewrote once or twice in the past few years, it's moving really fast. It's more of a standalone product for server metrics (node exporter + prometheus + grafana), wouldn't recommend to use as a general purpose database.
(Jaeger seems to have Kafka support so you can do this in real time. Haven't tried it.)