From Kafka to ZeroMQ for real-time log aggregation
tomasz.janczuk.org
tomasz.janczuk.org
I mean it's log/event aggregation for ops insight. Unless the whole system is some tightly coupled feedback loop into an unsupervised machine learning model where the whole thing has actual hard-real-time requirements (something which might well be impossible to build), then there's no possible way that having a second or two delay, between a message being created and when you can actually see it, can possibly matter.
I mean you don't have someone with instantaneous reflexes and resolution ability sitting there 24/7 with their eyes peeled as a stream of thousands of log messages flies by.
The whole premise seems spurious. Is it really necessary for every startup and their uncle to delude themselves into thinking that their use case is "mission critical" or "carrier grade" or whatever?
Auth0 provides basically "login as a service". Its not like they're managing the access control to nuclear launch codes or something.
Unless some medical device manufacturer was stupid enough to make a critical surgical assistance device require an internet connection, a WAN round-trip on unreliable networks, and reliance on a 3rd party service in order to start operating it... how can this service being down possibly be anything more than an annoyance? What's the worst possible scenario? A session has to be rebuilt? A user has to make an extra login attempt?
By their own admission the service has gone down already due to the old system architecture. How many babies died?
Also, the emphasis is not just the real time aspect. The article mentions the issues with kafka for HA
Live updated logs meant for human readability shouldn't need anything hard or soft real-time in the technical sense.
Zookeeper is a troublesome piece of software, especially when running on highly loaded clusters. I had experienced some pretty weird stuff with kafka + zookeeper, where reconnection and rebalancing of topics could break availability quite dramatically.
Another thing, that was not mentioned in the article is that Kafka guarantees the order of the messages, while ZeroMQ doesn't. This is certainly a showstopper for event sourcing applications.
On the other hand, ZeroMQ would be happily dropping messages when the pipeline is not balanced or when there is a spike in the data load. Also better set some sane high water mark, or you will end up consuming all your memory or probably get killed by OOM killer.
Anyway, while both are good solutions be careful not to be mislead by the advertising tone of both projects. Neither of them is a silver bullet...
Per Zookeeper, I'm sure it's better than each system trying to rebuild consensus on their own, but I am loathe to ever be responsible for running a cluster.
I'm surprised Amazon doesn't have Zookeeper-as-a-service, or perhaps even better, a Zookeeper-as-a-service facade that actually uses Dynamo behind the scenes.
I wonder if this is something inherent to Kinesis's design, or if Amazon will magically make it faster at some point in the future. Do you have a suspicion/indication either way?
We also have cases where we need many consumers reading a stream, so having that effectively limited to 5 concurrent consumers is a pretty tough limitation. Their solution to this is to just have consumers sleep on a failed poll, but that doesn't really help if you want to scale out your ingestion. You don't have that problem with Kafka.
I'd be interested in your results, if you were shooting for a specific number and couldn't hit it, just for future reference.
Single ordered delivery is hard. Really hard. It's easier to allow multiple delivery and/or unordered delivery -- if your log entries are given various correlating UUIDs and signed, this makes the whole problem about fifty hojillion times easier.
Loggregator[1] and Doppler[2] use UUIDs to identify request, app and agent, which makes it easier to allow logs to arrive out of order from multiple sources and then be reassembled.
One thing that helps a lot is to separate metrics from logs. Logging frameworks get overloaded into metrics systems. Logs are useful if you are recording essentially unique events ("I started at ...", "there was a request for /foo at..."). They are wasteful of high QOS resources if you're looking at statistical data in which fine-grained event identity isn't relevant ("57Mb is used by this process", "this request took 100ms").
For metrics it is better to use lossy sampling and derive a statistical view of goings on, rather than trying to capture every single data point.
A log entry should be like a page in your personal diary on your birthday. A metric is a phone survey asking your age. If you don't answer the survey, the data is still useful even with standard error.
Edit: I forgot my usual disclaimer that I work for Pivotal Labs, a division of Pivotal, which is the leading contributor of engineering effort on Cloud Foundry.
[1] https://github.com/cloudfoundry/loggregator
[2] https://github.com/cloudfoundry/loggregator/tree/develop/src...
Since we've decided to scope out access to historical logs
from the problem we were trying to solve and focus only on real-time log
consolidation, that feature of Kafka became an unnecessary penalty without
providing any benefits.When the system is overloaded.
Exactly the type of situations where you don't want to slow the system down further by spending resources on trying to let non-essential services survive.
Presumably their tradeoff is that it's more important for the system to remain available than for every log message to be delivered. Then secondly you try to deliver log messages with as high reliability as possible.
Often it is better to design for non-essentially components to fail early, or at least prevent their resource usage from escalating and dragging down other parts of the system (in this case, fixed buffers in 0MQ lets them isolate load in one part of the system by simply locking the rate the drain the buffers at below a suitable threshold that's normally fast enough).
First, message rates with ZeroMQ are often hundreds of thousands per second. The architecture must be designed so that no buffers, anywhere, overflow. If they do, you have a problem, usually a slow subscriber. Throwing out older data doesn't cure the problem. What ZeroMQ does is punish the slow subscriber by dropping so that the publisher doesn't crash. It's not recovery for the subscriber, it's protection for the publisher (and thus for other subscribers).
Second, trying to delete old messages is complex and sometimes impossible (if they're already in system buffers). The design of ZeroMQ's internal pipes has one writer and one reader, without locks. For the writer to mess with the reader would slow down everything and introduce risk of bugs. Dropping new incoming data is the only way anyone has ever found to keep things running at full speed.
These design choices were often delicate and counter-intuitive, yet they have turned out to be mostly accurate.
I'm not familiar with ZeroMQ's data structures, so forgive my ignorance. At the high water mark, why can't the consumer throw away old messages instead of the producer throwing away new messages? There are no locks or bugs — that's what the consumer does anyway.
I'm not saying that should be the default behavior, but perhaps an option. It's cleaner than silently dropping new messages, then sending an entire buffer of old messages if the subscriber recovers.
Also, as you've alluded to, guaranteeing log availability provides an enticing cross-tenant denial of service vector.
Vomit enough log messages into a shared fabric and you can begin to affect your neighbours. ZeroMQ, if I read right, diminishes some of this risk.
At Auth0, high availability is the high order bit.
We don't stress as much about throughput or performance
as we do about high availability. We don't have a
concept of a maintenance period. We need to be available
24/7, year round. Given the nature of some of our
customers, if we are down, "babies will die".
Though, as you point out, they simply dropped durability from their requirements.Then again, I've never tried running stateful services on Docker. Seems like a bad time.
> Kafka/Zookeeper combo struggled to rejoin the cluster
We've had issues with the broker properly re-registering with Zookeeper, but nothing that wasn't solved by stopping the broker and starting it again after the Zookeeper session timeout elapsed.
> Large companies like Netflix may be able to spend some serious engineering resources to address the problem or at least get it under control.
We're a company of 70, of which maybe 15 qualify as backend engineers/ops. Never had an issue managing Kafka. We've accidentally deployed Kafka with insufficient heap space and it kept on ticking.
> As a result of this difference, despite Kafka being known to be super fast compared to other message brokers (e.g. RabbitMQ), it is necessarily slower than ZeroMQ given the need to go to disk and back.
Writes do, certainly. But you don't have to "go back" from disk with Kafka during steady-state operation. Kafka writes segments to disk, this ends up in the pagecache. Reads to this segment are served directly from RAM which is excellent for the fanout consumer case. A healthy Kafka cluster usually has little disk read IO (unless a consumer is catching up or batch jobs are reading older data), but lots of network IO out.
Not saying that Kafka was the right solution here, but it seems like guaranteed delivery and archival of logs is vastly more important than "real time" logs, but what do I know? Their design of the Kafka-based system seems sketchy: colocating brokers with workers is awfully strange.
If logging went down, do babies die? If so, ZMQ seems like the wrong solution. If they just wanted to publish logs in realtime only, cool. I'm sure ZMQ will serve them well as it's a better fit for realtime non-durable publishing of data.
That said, I'd be interested in them publishing more details about their Kafka outages. The fact that we've had such contrasting experiences points to an interesting X factor in their setup that should be avoided by others.
He later rewrote the whole thing in C:
OSes and popular OSS libraries have ridiculous bugs sometimes.
It's less a silver bullet, more really bad tasting werewolf repellent you have to take every 2 hours.
I should clarify syslog, with naive configuration, is slower than file I/O because it does stuff, then writes to disk. feel free to correct me if I'm wrong
syslog-ng (my syslog of choice) has 3 'syslog' network protocols. `tcp`/`udp`, `network`, and `syslog`. Now lets play 'match the syslog-ng name to the rfc` (or lack thereof)!
HTH,
Regards, Robert syslog-ng documentation maintainer
If you want to do throughput of 10k/s per machine, do yourself a favor and use ZeroMQ. Kafka is a nightmare when your topology needs to change. ZMQ is just connecting pipes and splitting throughput.
Sending them from themselves to themselves you can do even better than that. Loopbacks aren't useful metrics for message passing. http://bravenewgeek.com/tag/rabbitmq/ is the most recent attempt to benchmark the various solutions (that I have found) and my own testing results in slightly less than the results shown...which I can only attest to unfamiliarity with RMQ, Kafka, etc. In the end, ease of maintenance only speaks to the weaknesses of solutions that need tweaking for specific topologies.
Somewhat off topic but you aren't worried about using docker to sandbox user scripting? Especially being a security company.
http://mikeangel1.deviantart.com/art/Metamorphosis-Franz-Kaf...
They're probably trying to optimise for making developers happy.