Distributed Logging Architecture in the Container Era
blog.treasuredata.com
blog.treasuredata.com
To solve this problem, you need a unique trace ID and vector clock, and this actually needs to propagated through your logging system, not through container tags.
Imagine you have 3 services. A calls B calls C. A request goes into A. The request fails. Since there are 3 services involved, you want to be able to see the log statements of all 3 services, in order. To do this, you need a UUID that is passed from A to B to C and a vector clock that is incremented from A to B to C so you get perfect ordering. Passing the data around is not super hard (you can use HTTP headers, for example). The tricky bit is having a library that does the right thing in A and B and C. And if A is written in Java and B is written in NodeJS, that library better work on both!
Just by adding a UUID to the requests as they're initially received I'm able to follow the chain of events across 5 different services, a message queue, and reconcile with any code exception reports from Bugsnag.
We don't have a lot of unique clients, but we do have a lot of requests.
Loggly can be set to parse logs (or have the data sent directly to it) and then tag messages with keys/values that can let you watch your request through every server.
Also, if I'm looking at access logs I'm not going to have that ID unless I start putting it in request headers. Right now the application logs and other monitoring tools pick it up though.
Edit: Didn't see mention of the clock on my first post, but it still stands for us that we're not in the millisecond timings between requests going in and out of services so we haven't had a problem with the timestamps from the systems they're on. Requests generally take a few seconds to process (handle a file upload of decent size, ETL some data, pump out some output document, etc) so each step through a service buys enough time before moving on to the next thing that we're not getting out of order log statements across each service.
Vector clocks don't necessarily have much to do with clocks. It's an algorithm for generating order in a distributed system.
The key is not just to have a single request ID, but also to provide hop IDs. Relying on clocks, even vector clocks, isn't ideal in a share-nothing architecture.
There's support for a Spark-job that generates a visual map of all your services' dependencies[0] where the strength of their relationship is represented by the thickness of the lines.
But that doesn't really have anything to do with the logging architecture. How long the thing lives doesn't matter when you're using aggregation, bypassing the local filesystem, and inspecting through a central portal.
People have been doing this in VM environments for a lot longer than they've been using containers. Whole companies and services were born to facilitate this, and some have even died already, in the time before containers hit full stride in the hype cycle.
We've been building distributed systems (with containers and other approaches) for years and this is the first time I've heard of logging being a problem.
It's not pretty -- there's no real-time alerts, but I'm not paying dollars per gb to AWS and I have no clue what google will end up charging... and I'm already running mysql. So there is that...
I've done a similar thing with Mongo in the past with their capped collections.
In particular, the ability to live-stream ("tail") logs seems to be a feature generally missing from logging aggregators. There was one (Loggly?) which provided a CLI tool for tailing, but it wasn't very impressive. (All the SaaS apps I have looked at do quite poorly when it comes to rendering live logs in a browser, too.)
Any recommendations?
[0] https://www.sumologic.com/press/2016-01-21/sumo-logic-announ...
(*disclaimer: I'm one of the co-founders)
However, the pricing looks very silly. 10GB/day is toy volumes, and only 30 days' retention, and this for $400/mo? That's not going to fly (sorry again).
Yes - screens are from old UI. New UI is currently in public beta and, as you'll see, dramatically improved:
https://www.scalyr.com/product/new-ui/opt-in
Re: pricing. A much bigger topic than there's room for here on HN, but in short - for a service that can ingest 1TB/day+ of your logs and gives you search times measured in milliseconds, you'll actually find our pricing is not only in line, but below a lot of other commercial providers (Splunk, Sumo Logic, etc.) And when you compare TCO with open source solutions (ELK, etc.) that require a significant amount of effort to scale, it's similarly competitive.
I suspect you're going to feel some competition from Google StackDriver/BigQuery. You can load 300GB/mo (your biggest plan) into BigQuery for about $156/mo including storage.
As someone who is basing our stack on Kubernetes, "effort to scale" is pretty damn small these days!
It's self-hosted, not a service, but it is rock solid, and it can easily sustain tens of thousands of messages being sent to it per second. Docker can send Graylog logging messages natively.
Use it. I promise you won't be disappointed.
If you'd like, I can send an initial email to you (from your HN profile) and you can bounce questions off me if you'd like.
And yes, that would be very nice of you!
So, new stuff goes to elastic, gets indexed, you can look at it via Kibana, or build custom dashboards, even directly from the logstash firehose.
As logs get older, you can delete whole daily indexes from ES, and if you want to investigate/datamine/aggregate something, you can still grep the archived logs.
The bottleneck will be probably Kibana (or the admin/operator looking at the end result), as all the other components can be scaled (beats are already per-node, logstash is stateless, so just run more of them behind a round-robin DNS name and beats will pick one up - or of course you can use HAproxy to load-balance, and the elasticsearch cluster can be rather large too).
If you are talking about tailing log files live, Fluentd has supported it from Day 1: http://docs.fluentd.org/articles/in_tail
Also, as other sibling comments mention, there are tools, both SaaS and open source, that you can use as a destination of the logs Fluentd tails/listens/collects.
* Elasticsearch: https://www.digitalocean.com/community/tutorials/elasticsear...
* Graylog: http://www.fluentd.org/guides/recipes/graylog2
* Scalyr: https://github.com/scalyr/scalyr-fluentd
* Loggly: https://www.loggly.com/blog/stream-filtering-loggly-fluentd/
* SumoLogic: https://gist.github.com/d-smith/8d3e7d53db772c6a7845
* Papertrail: https://github.com/docebo/fluent-plugin-remote-syslog
(and literally hundreds of others)
>though I rather wished Heka had taken off; it's much more flexible and in theory leaner and faster since it's Go
Heka was a great project, and a drop-in binary (as opposed to requiring a VM like Ruby) approach was interesting if not compelling in certain situations. That said, I never saw any benchmark that showed Heka was materially faster than Fluentd (or Logstash, for that matter). A lot of speed in this type of complex software comes from data structures/algorithms, an appropriate use of low-level language bindings, etc.
While language plays a role in the speed of software, it's hardly the only factor. As you said, it's only in theory, not in practice =)
For example, you will definitely want to filter on labels (including regexps) while tailing, and such filtrering should support adjustable context (both # of lines and time interval) and should support an optional time range to scroll back into history ("tail from 2pm"). And of course, grepping of historical entries.
Another thing I am not impressed with is pricing. Loggly seems the most reasonable in terms of price per volume, except it limits the number of "team members" to a ridiculous degree (5 users or something like that).
I set up Graylog once and wasn't impressed with it. Its reliance on Elasticsearch means it is quite static when it comes to the schema/input format. You can't change the settings; there is no reindexing support (or at least this was the case when I tried it, a year ago or so). Also, don't think it has any CLI tools?
I was contemplating trying out logstash or fluentd, but I don't see any major advantage over current simple solution. Can anyone more knowledgable than me help with explaining what I would get from logstash/fluentd?
If you just write to stdout everything makes sense in the file top to bottom, but if you're logging json objects into a database I had problems; if your parent process echos stdout from the child you lose a lot of context about the child, you can also double-log child process messages (because the parent is also echoing them), and you can't easily associate the parent/child objects. I came up with a hacky solution, but I wasn't logging these events to a centralized server at the time.
(I haven't used gcloud's dashboard, so I might be making wrong assumptions)
To respond your question: it sounds like you have something that works fine for logging messages. Personally, I split it into categories 1) logs (serial events that needed context for any usefulness) and 2) events (metrics or exceptions). I wrote 1) to traditional log files and wrote 2) to an ElasticSearch or statsd database to log exceptions to. Ideally, I wanted to use the same mechanism for both and peel off relevant data into separate databases.
A metrics database like Elasticsearch will let you query things like, "What modules give the most errors?" "Has this function been called more often this week than last?" "Is this process taking longer when using the newly released version compared to the old one?" etc.
Whats working for us now is local docker json logs -> heka -> kafka -> graylog -> ES. We even dockerized graylog to scale the processing up and down on demand.
Does 250k messages/sec at peak easily, More details in this presentation https://www.youtube.com/watch?v=PB8dBnpaP8s
250k messages/sec is pretty low.
It'd be useful (not to mention entirely feasible) to be able to handle 5M/sec bursts on a single system and aggregate sustained 50M/s+ (10Gbps network).
The system should be able to establish total order (serialization) locally and causal order over the network.
I did some tests and found out hardware limits are somewhere between 20-50M messages per second (serialized) on current consumer grade X86. Practical implementation would of course be slower.
Of course you need to go binary logging at that point, adhere to cache line boundaries, etc. mechanical sympathy. Maybe even do usermode networking at the aggregator server.
Binary logging is a must, because even something like string formatting is simply way too slow, by an order of magnitude.
In my quick tests I found the string processing hit to be surprisingly high, 10-50x. C "sprintf" and C++ stringstream are atrociously slow. (Surprisingly, considering sprintf has a pretty complicated "bytecode" format specifier parsing loop, it was still significantly faster than stringstream implementation.)
Timestamps are another huge performance issue. Typical system calls for high precision (a few microseconds or better) timestamp take several microseconds each - using one of those will alone drop performance to 100-500k range. It's of course much faster to use RDTSC, but then you have the issue that different CPU sockets have different offsets (sometimes large) and possibly also some frequency difference between them. Also older CPUs don't have invariant TSC, so their speed changes when CPU frequency changes.
(Some word of warning about Windows QueryPerformanceCounter: When you test it on your development laptop, it appears fast, because it's using RDTSC behind the scenes. But when it runs on a NUMA server, it often changes behavior and becomes 20-100x slower, because Windows starts to use HPET instead. RDTSC takes maybe ~10 ns to execute and it's not a shared resource. HPET takes 1-2 microseconds to read and is a shared system resource, concurrent access from multiple cores will make it slower.)
> Binary logging is a must, because even something like string formatting is simply way too slow, by an order of magnitude.
true, string processing takes most of the resources.
Some solve it by caching a low res timestamp, but that means losing ability to get side channel information about precisely how long time passed between two events. High resolution (micro/nanosecond range) timestamps give so much more information than just time.
To give a rough point of comparison, CloudFlare is taking 4M page views per second and that's 5% of the global internet traffic.
No. Although I've done that too...
To debug timing critical code. Or Heisenbugs that disappear when you log. Ever had a bug that vanishes when you enable logging?
Not 5M sustained, just short sub-millisecond bursts here and there.
Back to the topic. I suppose that saying "5M/s" to refer to 5000/ms during a single millisecond is misleading.
There are many implicit "buffers" (raw CPU power, intermediate hard drives, TCP queues, syslog/fluentd processing) that may smooth that kind of peak very aggressively long before it reaches graylog.
Interesting talk by the way. I'm doing similar things, at a smaller scale though.
I'd be curious to know some more details on this. I guess Go channels do copying rather than sharing, but I'd still expect Heka to perform better.
Personally I've found Heka to be more robust - I've had to fix bugs in plugins for the former two and generally I haven't found them to be architecturally sound.
For example with logstash I've had bugs in plugins which would crash the entire daemon -- that's just inexcusable for a core infrastructure service.
The multi-tiered thing sounds very strange, why not just hold an on-disk buffer?
The need for a performant and powerful log shipper is still there. I hope to see some new options come around soon that can achieve 1MM+ lines/sec from a single daemon without requiring multiple tiers, receivers, etc.
It's pretty similar to Fluentd in architecture, some features are:
- Event-Driven (async network I/O).
- Input / Output plugins.
- Data routing based on Tags.
- Optional SSL/TLS for networking operations when required.
Next major version 0.9 will come with buffering support (memory/file system). Ah, it's fully made in C.
I see that your input/output plugins are written in C[0]. I'm guessing this is because of the constraints of the embedded environment, but it really doesn't seem like it would be worth it in a normal one. The LUA sandbox model (e.g. Heka) just seems highly preferable.
My main problem with Logstash/Fluentd is precisely the fragility and non-robustness of the plugin system.
[0] https://github.com/fluent/fluent-bit/blob/master/plugins/out...
The decision about "why C" is: flexibility, performance and adaptability (note that it was originally designed for Embedded Linux targets, but now going everywhere). In order to make things easier for output plugins, every time a set of records needs to be flushed through some output plugin, a co-routine is created so any plugin can yield/resume at any time. For example out_http, out_es and out_forward relies on network I/O, having an event loop and a coroutine associated allows to simplify the plugin development and state management: connect, write, read, etc. This model is the foundation and allow the next step to integrate scripting more smoothly. For environments without co-routines support (old compilers), a POSIX thread model exists.
What are the specific "fragility"/concerns you see in Fluentd plugin model?
E.g. Logback is perfectly happy shipping structured logs directly to Kafka or Elasticsearch, with no need to re-parse the formatted log output.
It's a well-understood pattern that Heroku made more visible under "The 12 Factor App".
stdout/stderr require no configuration in any language that I'm aware of. Why should the app care about how to wire up logging? That's a platform concern.
On Cloud Foundry you get Loggregator, which frankly needs improvement, but for the most part you don't care about how to wire up logging. You print to stdout or stderr and the platform wicks that away to a firehose service for you. You can hook up kafka to spout, or elasticsearch, or anything else you like.
Disclosure: I work at Pivotal, we donate the majority of engineering to Cloud Foundry.
Furthermore, I've got my graylog setup running from a customized docker-compose.yml file, and this thing sings.
Also, never had network contention because I logged too much unless something was really broken - and then only at that node.
I'll continue sending events to my anycast address without any aggregation thank you very much.
> They say that without aggregation scale-out is
> impossible. That's simply not true. Using things like
> anycast and ECMP it's super easy . . .
Then you haven't worked at significant scale. That's fine! Just keep everything in context :)I guess it boils down to what you call aggregation.
Is having multiple stateless receivers behind anycast/ECMP aggregation writing to a distributed database aggregation? I'd argue not, but maybe this is where the difference in opinion lies.
http://logz.io/blog/fluentd-logstash/
https://www.pandastrike.com/posts/20150807-fluentd-vs-logsta...
That said, I use nxlog basically everywhere because I run a heterogeneous environment and it works well on all OS's I use (and is fast and light on resources)
syslogd(8) -r This option will enable the facility to receive message from the network using an internet domain socket with the syslog service (see services(5)).
Why won't this work with dynamic DNS ?