Eliminating Task Processing Outages by Replacing RabbitMQ with Apache Kafka
doordash.engineering
doordash.engineering
Option 1: Redis broker. Did they ever setup a test environment with a Redis broker to see how it compares to a RabbitMQ broker? Seems like they didn’t.
Option 2: Same as option 1. Did they actually try to setup a kafka broker for celery to test it out and see how it goes?
Option 3: Multiple brokers. They could have attempted a very simple setup starting with just two RabbitMQs and dividing the tasks roughly in half between them. This would allow scaling of RabbitMQ horizontally. At least try it and see.
Option 4: Upgrade versions. Did they even try setting up an environment with upgraded versions. They say there’s no guarantee it fixes their observed bugs. How about trying it out and doing some more observation?
Option 5: We decided to go all-in on this one without even trying the other 4 and doing any sort of comparison with actual tests or benchmarks. Maybe this was the best solution for them, and they chose wisely, but how do they know?
The big thing that gets me is they talk about the observed issue of celery workers that stop processing tasks as well as the limited observability of celery workers and RabbitMQ. These aren’t black boxes. They are open source. You can debug them yourselves, report the bugs, fix the bugs and submit patches, add features, fork it yourself if necessary, etc. The fact that they don’t know if new versions will fix their observed bugs makes it clear they never identified those bugs. What’s the point of running on an open source stack if you’re going to treat things like a black box anyway?
The most blatant example is the countdown tasks. Celery has a very strange implementation of these (meant to be broker agnostic) where it consumes the task from queue, sees in the task custom headers (which is meaningless to RabbitMQ) that it should be delayed and then sits on the task and takes a new task. That results in heavy memory load on your celery client holding all these tasks in memory, and if you have acks_late set, RabbitMQ will be sitting on many tasks that are claimed by a client but not acked and _also_ have to sit in memory. But that is 100% a celery problem, not Rabbit, and we solved it by overriding countdowns to use DLX queues instead so that we could use Rabbit-native features. Not surprisingly, Rabbit performs a lot better when you’re using native built-in features.
I, too, would love an implementation of arbitrary-granularity delays without buffering (potentially huge, mem-wise) tasks in "unacknowledged" in a giant binary heap: that's expensive for consumers and terrifying for RabbitMQ stability when a broker has to e.g. re-"ready" millions of messages because a delay-buffer OOMed.
https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages...
The plugin keeps a node-specific database of the messages that are to be delayed. If the node is unavailable or lost, so too are the messages the node was keeping for future publish.
Also, there is no visibility in to the messages that are awaiting delay. 10? 100? 5000? Your guess is as good as mine and you'll never be able to figure out what's in to-be-published pipeline.
Signed, someone who has dealt with these issues.
The plugin keeps a node-specific database of the messages that are to be delayed. If the node is unavailable or lost, so too are the messages the node was keeping for future publish.
Also, there is no visibility in to the messages that are awaiting delay. 10? 100? 5000? Your guess is as good as mine and you'll never be able to figure out what's in to-be-published pipeline.
Signed, someone who has dealt with these issues.
This will probably not be a popular opinion here but the business of DoorDash is not to fix open source software.
I get your point. I use a lot of OSS in my business and we do contribute back. But not every business has the opportunity to observe the stack for days or weeks while the real clients suffer due to outages. Each client who is affected by an outage is most likely churning immediately and you're not going to win them back easily.
In essence, the migration to Kafka will only benefit them long term and they have ultimately made the right business choice.
It probably isn't. But as their business depends on software, it's worrisome that rather than understanding what was wrong, they just switched software. That doesn't seem inspiring to me.
Nobody said they had to fix the problems, but identifying them, and checking if they were fixed upstream seems like a good idea. Sometimes upgrading is a nightmare and doesn't fix anything, sometimes it fixes a major issue and drops cpu usage from 90% to 5%. If you know what the problem areas are, it should take less than a day to look at the changes and see if that area was addressed. Even if not, reading the changelogs and version notes should give some idea. It's fine if they didn't test upgrading, but again, I'm not inspired about their engineering talent and processes when it looks like they didn't even look to see if it would have been likely too help.
Engineering Daily podcast has various podcasts that come up around this. One that comes to mind was the one with Slack. Slack had hyper-growth, they went out and hired individuals who had already tackled these problems at the same numbers or larger.
You get hired at a company, and you look around at what FAANG is using, you implement that because you can always say, "But it works at Netflix". You build up immense expertise in the project, the tool, etc.
Switch jobs, ignore everything the existing company is doing, where the landscape is going and shove your agenda of implementing the same thing such as Kafka. Measure everything.
Look like the saviour of the company. Point to this, ask for massive stock compensation, bonus, etc. based on the metrics that make you look good.
1. Cache it's amazing... unless you need keys something more complex than key value, e.g. multiple keys, invalidate by pattern etc... So ... not so good actually.
2. Pubsub doesn't sale really good, and that'a a real limitation
3. The new pubsub is purely a copy of Kafka ( they admit that in docs) but in a much less supported, and feature rich version.
So despite I'm still using Redis in my Day today, as of 2020 I don't see redis in any new setup.
RabbitMQ has very good observability. It has built-in web dashboards as well as an API. Every day that I had to bang my head against the opacity of an org-wide Kafka or ActiveMQ setup at one job, I missed the simplicity and transparency of RabbitMQ from a previous job. Somehow Kafka earned the "trendy" badge though, so everyone just uses it by default.
This is unfortunate if you want to see the queues from up close.
I don't have the full picture, but I'd move to multiple brokers first, which seems like a low effort move, then away from Celery, in order to split worker code from the rest of the application. In the meantime, I'd instrument the hell out of Celery to see what was going on and push those changes back into upstream.
But the first stack I would reach for to handle a queue would be erlang/otp (or maybe elixir), is that not designed to handle large queues without failure?
Again I am absolutely no expert but I'd like to hear the views on this.
In absence of any significant stats its safe to assume their decision were driver by this two statements from article.
"There were no in-house Celery or RabbitMQ experts at DoorDash who we could lean on to help devise a scaling strategy for this technology."
"DoorDash had in-house Kafka expertise"
Ah, yep, as usual you end up with a queue for your queue because Kafka isn't actually a queue in the first place. This is why I recommend Kafka for _data_, rabbit for _tasks_. The article doesn't really explain how this works, but I'm guessing that one service holds commands in memory, then sends http messages to the workers. As they noted, you end up with lost tasks on restart, and that's the big gaping hole that a lot of people would not want to fall into - your mileage may vary. I assume the kafka reader service has to also track how busy each worker is, also lost on restart.
In the bit you quoted, I'd assume the local work queue has callback objects.
> Failovers took more than 20 minutes to complete and would often get stuck requiring manual intervention. Messages were often lost in the process as well.
Data loss while RabbitMQ is down is certainly a problem.
> The Kafka-consumer process is responsible for fetching messages from Kafka, and placing them on a local queue that is read by the task-execution processes.
This is a bit confusing, so now instead of a centralized outage, grinding everything to a halt, the risk is distributed data loss as local queues go down? As long as they stay down, it should be limited data loss to the size of the queues times the number of consumers, but if an instance is flagging it could eat through a whole lot of messages before someone notices the problem. I guess it’s really dependent on the types of messages being processed and if they are idempotent enough to be replayed without consequence.
Kafka is not a queue. It's a distributed log. The difference is significant because in a queue, the message can be processed by a single worker within the single unit of time - message is picked, locked, nobody else can see it, after processing, message is either acked or nacked.
In a log, multiple consumers can process the same message in parallel. If a message is bad, it is usually skipped by simply advancing the offset to the next offset.
The log has the advantage that the reason for slow processing is easy to identify and isolate - it's the slow consumer. It's also very easy to monitor the consumption lag - simply look at how fast the maximum offset grows vs what's the current consumed offset (high / low watermark). If the difference grows, your consumer is too slow.
Like...the whole point of having messages work that way is that you're not having to bake event processing and duplicate transaction handling and compensating transactions and all the other message handling logic into literally every process which comes after your message queue.
It's a non-solution - you're just leaking the problem out to another part of the stack.
From the user perspective, one only needs to ensure having enough consumers in a consumer group.
At Klarrio, we have a rather small cluster in production but we can easily push 2.4Gbps with <100ms end to end latency for dozens of consumers.
It’s working really well.
> This problem was never root caused, though we suspect an issue in the Celery workers themselves and not RabbitMQ.
So the problem wasn't Rabbit itself, but the usage/lack of understanding.
It's a great writeup. But I'm personally curious about them hitting the limits of vertical scalability. They said they used the biggest node 'available to them' without clarifying how big that was.
If the issue was RabbitMQ hitting resource limits then upgrading its hardware (or swapping it out for some other similarly featureful broker but which is faster) might have fixed some of the problems too, albeit not things like insufficient observability.
I also wonder to what extent RabbitMQ suffers from being written in Erlang. Switching to a featureful broker with sharding and replication that isn't Kafka e.g. Artemis might have allowed them to avoid rewriting code to not use helpful features.
This performance blog suggests RabbitMQ may top out at ~4000 messages/sec when making things durable, which isn't especially good performance:
Would I personally reach for kafka out of the box? Not so much. I personally think it is quite a monster to get setup correctly. If you need dual datacenter up time it also becomes very expensive as you need to buy cloudera or confluent. Or have a plan to manage it yourself which will be 'interesting'.
MirrorMaker is pretty decent and part of the oss release. There are also multiple alternatives available. There is no need to buy anything.
I actually tweeted out some similar issues a week ago: https://twitter.com/jarshwah/status/1310820638655877120
There are ways to scale rabbitmq past 1 node without giving up too much throughput.
- use multiple queues as each queue is single threaded
- use sharding plugins which present a single queue externally but use multiple queues across a cluster internally
- use lazy queues to get messages straight to disk, lower throughput but higher message count
It doesn’t really sound like the other options were seriously considered and it was more about moving away from celery. That’s fine - just be honest about it.
However, that's less Celery's fault and more RabbitMQ's: pub confirms are a Rabbit-specific extension to their protocol, and are left off in most example implementations. MongoDB, also, chooses horrible default settings; however, I wouldn't blame a Mongo-backed product entirely for being hackable.
Additionally (I mentioned this in a sibling thread), Kafka's publish behavior is even worse. The vast majority of Kafka publishers' default settings don't even send the data to the broker without waiting for a response (like RabbitMQ without publish confirms); they send the data to a local buffer which is asynchronously flushed periodically/volumetrically. If your process crashes, you lose the buffer. Unpleasantly, if you force a buffer flush on every publish, Kafka's ingest rate goes from 10-100x RabbitMQ's to 0.01-0.1x RabbitMQ's.
If celery is going to go to the lengths it does to abstract brokers then it should also mention required configuration for reliability.
When you publish a message to Kafka, it stores the message in a buffer and immediately returns a future that completes once the data is sent to the broker.
If you forget to wait on completion of the future that is on you, this behaviour is well documented.
Absolutely. However, data loss and edge cases emerge where process crashes and abrupt exits are concerned. With RabbitMQ, even with pub confirms disabled, the odds are that all of your publishes (sans, perhaps, the one you were in the middle of at the instant of exit) will make it to the broker. With Kafka, every message not yet flushed from the buffer will be dropped.
I get that this is expected/documented, and I get why it exists (batchwise operations allow way higher net throughput into Kafka). However, I've encountered many groups of engineers from many companies who were surprised by and unable to effectively work around this behavior. Workarounds include performance sacrifices (flushing per-publish); complex signal handling semantics to coordinate flushes (complicated further by errors that occur if a signal arrives during flush); and atexit(3) flush hooks (terrifying, given the way some higher level programming languages' garbage collection interacts with atexit).
At the end of the day, it matters how friendly a tool is when it's held by a novice. While there are many good reasons to use Kafka, that's not one of them.
We also evaluated kafka vs. rabbitmq, and chose rabbitmq.
For people hitting performance degradation with rabbitmq this is a good reference [0]
RabbitMQ should be used as-is without a middleman as much as possible, to make the best possible use of all of its concepts.
[0] https://www.cloudamqp.com/blog/2018-01-19-part4-rabbitmq-13-...
If you want redundancy or write-scaling, you can run 2+ of them side by side, round-robin your jobs to the available instances, and have your workers read from them (RR as well).
Beanstalkd is a lifesaver for job queues. It's the first thing I reach for when I need a queue. I'm excited to see how Redis' Disque turns out, but until then I'm happy with bean.
I've evaluated replacing bean with something like Kafka, but it just doesn't make any sense at all. A distributed log is NOT a queue, and even using something like RabbitMQ gets tricky if you want to pull jobs instead of push them. For dedicated worker queues at small-to-medium scale, I cannot recommend beanstalkd enough. No, it doesn't have native sharding or multi-region replication, but like I said, I run millions of jobs/hour through it and never needed any of that crap.
now that rabbitmq clustering is actually working (~7y ago it was basically impossible to setup clustering without having huge headaches), you can easily setup 3~5 machines in cluster and never have a single issue.
How this problem of scheduled /delayed tasks was solved after moving to Kafka? The another system they mentioned, is something else entirely than the solution proposed in the post?
Problems with observability? We added lots of custom StatsD and text logging instrumentation (Celery "signal" middleware), so that we could get e.g. accurate "how long do tasks like this spend waiting in the queue?" answers. Other than the inherent limitation of RabbitMQ being a strict queue (unlike Kafka, you can't "peek" at things in the middle of a topic without consuming that topic--we could do something hyper-complicated with exchanges and deliberate duplication of messages to address this limitation, but that doesn't sound worth it to me at all), observability of the brokers themselves seems pretty good. Coupled with Sentry reporting issues that occur during task processing, and some custom decorators to retry/delay/store tasks affected by common classes of systems issues, our visibility tends to be better than I've seen in any other asynchronous queue/message-bus system I've worked on.
Sneaky task disappearances turned out to be mostly bugs in the way we were starting/stopping workers, related to signal handling and rabbitmq "redelivery". By really understanding the Celery worker start/stop lifecycle and coupling that with how Systemd chooses to kill processes, we were able to reduce those to zero. Celery also had a couple of "bugs" (questionable behavior choices) in this area which were resolved in 4.4.
Celery ETA/countown task induced RabbitMQ load turned out to be because Celery made the questionable decision to queue ETA tasks on every single (eventual executor destination) worker node. We customized the celery task-dispatching code to route all ETA tasks to a set of workers which only buffer tasks, and re-deliver them to the actual work queues when their time is up. As a result, the domain of cases in which RabbitMQ had to "take back" (unack -> ready) large amounts of tasks went from "every time every worker restarts (and we deploy to them a lot!)" to "every time a very specific, single-purpose worker crashes", which reduced issues with that system to zero.
A lack of scale-out in RabbitMQ was addressed by adding additional separate brokers (and therefore celery "apps") along two axes: sometimes we peel off specific task workloads to their own broker, and in other cases we run "fleets" of identical brokers that tasks round-robin across, with a (currently manual, moving in the direction of automatic) circuit breaker to take a broker out of rotation if it ever has issues. We wrote a whole blog post[1] about scaling that specific set of RabbitMQs and workers. Totally agree with Doordash that rabbit's "HA" generally isn't, and that scale-out needs to happen across brokers.
Connection churn-related broker instability was addressed partially by scaling the number of brokers, but also on the consumer side, by carefully tuning per-worker concurrency to minimize connection counts while doing as much work as possible on a given piece of hardware, and by disabling the Celery remote control channel (pidbox queues). While that means that nice tools like e.g. Flower aren't as useful to us, it also means that the cost to a RabbitMQ broker of losing a whole bunch of consumer connections is much lower. When it comes to the "harakiri" churn of publisher connections discussed in the article, we haven't encountered connection issues due to our publisher tier. Doordash's web tier is almost certainly bigger than ours, but I'd make a deeply uneducated guess that we're at most an order of magnitude apart. I'd be curious to learn more about the story there, since, even at a reduced size, we regularly run 10ks of connections on a single broker with a pretty high flap-rate due to e.g. recycling webserver processes or restarting consumers due to code releases.
In general, I agree with the article and yesterday's Celery 5.0 release post comments, that Celery is a quirky piece of tech, and that RabbitMQ is far from simple to run. However, I'm generally pretty pleased with Klaviyo's approach to go "through" the problem by diving deep on issues we had and fixing them in the stack we chose, rather than tossing large parts of it and re-learning the foibles of some other piece of software. At present, we run dozens of brokers and process 100ks of tasks per second at peak volume. While nobody considers our setup simple or issue-free, it's one of the most fully understood pieces of technology we run.
While it's not out of the question for us to adopt it at some point in the future, Kafka was discouraging to us when hardening our RabbitMQ/Celery setup for a few reasons (though we do use it for some other pieces of our infrastructure which require it):
As the Doordash folks indicated in the article, Kafka is really not well-integrated with the Celery stack at all, so building in things like front-vs-back-of-queue retries (both of which are extremely useful in different situations), deferred delivery, and the ability to rapidly change the number of consumers on a topic all take effort. Each of those problems has a solution, or at least a response, in the Kafka ecosystem, but Python task-processing frameworks which integrate those behaviors are both unfamiliar to us, and significantly younger than Celery.
As with any clustered (rather than sharded) system, we lack expertise in understanding why publishes fail when Kafka is in a partially degraded state. With our existing Kafka workloads, many failure waves (consume or publish) happen without a full grasp of what's wrong/how to fix it. That's most definitely an "us" problem, and we are learning, but it's likely going to be quite awhile before we're at the comfort level that we currently have with, say, recovering message data from a data volume in the aftermath of a massive AWS-induced broker crash, or replacing RabbitMQ nodes that are experiencing elevated latency or flow control.
Lastly, we were unpleasantly surprised by Kafka publishers habit of lying to their clients and saying a message was published when it was in fact buffered in-memory, pending a periodic (or volume-initiated) flush operation. Our processes crash a lot, usually when we least want them to, and having those crashes cause the data loss of an entire pre-publish Kafka buffer has been extremely unpleasant for us. When we reduced "batch.size" to 1, we discovered, to our dismay, that Kafka's vaunted "way better than RabbitMQ" publish time and volume numbers were entirely dependent on batch-wise optimizations, and that publishes were tens of times slower with batch.size=1 Kafka than they were with pub-confirms-enabled RabbitMQ (RabbitMQ also has batch-wise optimizations with publish confirms, which I'd argue have vastly better reliability semantics than Kafka, but that's another story and we're not using those...yet; ask me if interested). Again, that's partly an "us" problem (our crash rate is high, and dropped Kafka-destined batches could be recovered in other ways if we spent the time), but one that we don't have to worry about in the Celery/RabbitMQ setup.
1. https://klaviyo.tech/load-testing-our-event-pipeline-2019-42...
Edits: a few for clarity and removed a few things that the Doordash folks already covered and that my initial less-than-careful read didn't catch. The substance of my post didn't change.
The way I saw it, a message queue is a really fundamental piece of reliable distributed systems. It's fundamentally stateful and I would like it to never go down and certainly never result in dropped messages/tasks. Like with an RDBMS, if I can pay someone else to solve that problem for a reasonable fee, I will rush to give them my money. The next time I built a queued task management system, I used SQS. In the past 7 years I've used it, it's generally worked like a dream.
I can understand if for performance, cost optimization, transparency or portability reasons you still want to run your own RabbitMQ cluster. But my default bias is against that. So with that said, I'd like to ask a follow up question to this excellent comment: did you consider a managed queue like SQS? If so, why did you elect not to go with it? SQS even has a celery broker now (although certainly not as cleanly integrated as Redis or RabbitMQ).
Edit: additionally, there was an inertial component to that decision. Not only do some folks have pretty deep expertise with RabbitMQ, but by the time we started to hit significant limitations with it, we had devised pretty good (manual and automated) ways around those limitations, like the multi-broker approach. I can't predict what would have happened had we jumped to a totally different messaging technology when we ran up against those things, but I'm reasonably happy with where we ended up.
If you have folks with deep expertise with RabbitMQ, I can see why you'd go with it. It means that you can essentially neutralize its operational risks and maximize its strength, which is performance (as well as great integration codebase maturity).
> Kafka publishers habit of lying to their clients
I have no Kafka expertise but it feels like there ought to be a Kafka configuration option somewhere that tells it not to do that? (And in fact changing the batch size to 1 won't help if it continues ack'ing messages without syncing to disk)
Sorry, I may not have been clear. For many Kafka clients, setting the batchsize to 1/interval to 0 causes every publish operation to block on the batch being flushed internally. Flushes are acknowledged by Kafka (and, like any RPC, can succeed or fail), and their success indicates that the persistence system was engaged, indicating disk persistence to a configured number of disks.
Other Kafka clients provide a manual "flush" operation which acts similarly.
More information about, and discussion of this phenomenon in the Python Kafka client can be found here: https://github.com/confluentinc/confluent-kafka-python/issue...
The deceptive behavior only happens when a publish operation does not entail a flush. At that point, unsuspecting client code may have assume that data left over the network when it did not.
RabbitMQ's behavior with publisher confirms disabled is similar but not identical. Even with pub confirms off, most RabbitMQ clients "publish" operations are synchronous (as in they return after data has been written out over a socket). Publisher confirms are an added layer of resilience that enables clients to listen to RabbitMQ saying "I have successfully routed (and, depending on settings, persisted to disk) these messages". That acknowledgement may never come (rabbit may be overloaded, crashing, or may reject a publish), but on most happy-path systems data loss is not severe even without publisher confirms--turning them on is there to fill the non-happy-path case, and also to provide a very primitive system of backpressure to publishers, forcing them to wait on an overloaded broker rather than sending it even more data when e.g. its persister is slowing down.
Because of that behavior, I'd argue that even without publisher confirms, RabbitMQ's publish behavior is still less data-loss prone than Kafka's. Regardless, the only way to run either system in truly reliable "publish means my data is on the broker" mode is to set the bufsize/flush interval to the minimums (in Kafka's case) or to wait for a confirmation after every single publish (in Rabbit's).
As with all strategies to maximize reliability, those both come with a performance cost, and shouldn't be blindly adopted unless you have a good understanding of how much data and performance loss is acceptable.
If you can relax your constraints, then something like NSQ can be a great option that I heartily recommend.
edit: in this sense - https://en.wiktionary.org/wiki/Kafkaesque