223 karma · joined June 12, 2009
If you have a Faust app that depends on a particular feature we strongly suggest you submit it as an integration test for us to run.
Hopefully some day Celery will be able to take the same approach, but running cloud servers cost money that the project does not have.
It's important to remember that users had difficulty understanding the concepts behind Celery as well, perhaps it's more approachable now that you're used to it.
Using an asynchronous iterator for processing events enables us to maintain state. It's no longer just a callback that handles a single event, you could do things like "read 10 events at a time", or "wait for events on two separate topics and join them".
- A Stream iterates over a channel
- A Channel implements: `channel.send` and `channel.__aiter__`.
- A topic is a "named" channel backed by a Kafka topic
- Further the topic is backed by a Transport
- Transport is very Kafka specific
To implement support for AMQP the idea is you only need to implement a custom channel.
If you open an issue we can consider how to best implement it.
Simplicity is of course a goal, but this may mean we have to sacrifice some features when Redis Streams is used as a backend.
I'm glad you like Celery, this project in many ways realize what I wanted it to be.
Thanks for pointing this out. Fixed the links. You can also find the docs here: http://faust.readthedocs.io/en/latest/
Faust uses Kafka for message passing. The new messages you create can be pushed to a new topic and you could have another agent consuming from this new topic. Check out the word count example here: https://github.com/robinhood/faust/blob/9fc9af9f213b75159a54...
Also note that the Table is persisted in a log compacted Kafka topic. This means, we are able to recover the state of the table in the case of a failure. However, you can always write to any other datastore while processing a stream within an agent. We do have some services that process streams and storing state in Redis and Postgres.
Yes, we're definitely interested in supporting Redis Streams!
Faust is designed to support different brokers, but Kafka is the only implementation as of now.
This answer sums up why we wanted to use Python for this project: it's the most popular language for data science and people can learn it quickly.
Performance is not really a problem either, with Python 3 and asyncio we can process tens of thousands of events/s. I have seen 50k events/s, and there are still many optimizations that can be made.
That said, Faust is not just for data science. We use it to write backend services that serve WebSockets and HTTP from the same Faust worker instances that process the stream.
Disclaimer: I'm a contributor
I have merged many features, like broker transports, result backends, etc, and while the initial contribution was great, it ends up being unmaintained with issues that nobody fixes.
If there's any feature that you really want back, chances are the problems with that feature are not super difficult to fix, so please reach out!
It doesn't make much sense to me, unless they are carefully wording this into something that seems reasonable to the public ("I'll just use private mode, no deal") when it really means monitoring all our internet traffic. But then why is the media repeating it?
I only wanted to clarify the RAM situation described in your slides, as it's not widely known that you can use eventlet/gevent with Celery.
I want to clarify something about Celery and RAM usage.
When writing web crawlers, and other (mostly) I/O-bound tasks, you should be using the eventlet/gevent execution pools instead of the multiprocessing one. This will drastically reduce memory use, and perform better.
If you have four CPU cores you can start four worker instances with 1000 threads each (for a total of 4k threads): `celery multi start 4 -A proj -P gevent -c 1000`
This will utilize all the CPU/cores in your system, working around the GIL.
One of the new features coming in Celery 4 is a message protocol with support for multiple languages, maybe we could have an Elixir worker soon.
A greeting can be made from across the street, and sometimes just lifting your hat in response is enough (I don't have one, but I have seen people do this).
But that's my opinion and I have come to learn not everybody value succinctness in code, so you have the choice of both :o)
Had I not been constrained by backward compatibility I may have made it so that task(arg) only defines the signature of a task invocation and you'd need to do task(arg).delay() to call it remotely, and task(args)() to call it as a function locally.
I regretted this as soon as I submitted it. I would hate for someone to do the same thing to my projects so I should know better. I've written about it before, but realize that you probably have not read it :)
I really like the multiprocessing library, it helped me start Celery in the first place. What it tries to solve is actually very very complicated, and you would have to test it on production systems for years to be sure it works, and I think Celery was the app that did that testing. I contributed some fixes back into Python, but most of it is not merged upstream.
The most complicated issue I had to solve was that multiprocessing.Pool uses POSIX semaphores to share pipes between processes (that's how the pool processes receive jobs, and the parent receive results). If a child process is killed before releasing that semaphore you have a deadlock that's tricky, if not impossible to solve. So I rewrote the pool to use async I/O instead, which also had the side effect of drastically improving performance (no locks). Sadly I'm not sure how to implement that on Windows, so it's unlikely to be merged upstream. Other fixes and features used by Celery is available in our billiard (on PyPI) fork of multiprocessing, but the async pool is not part of that yet as it currently depends on code in celery that does not fit in billiard (it should be rewritten to use asyncio now).
You can claim to replace Celery using a small layer on top of async I/O, or claim to replace Celery with a simple Redis list operation, but I think that's unfair to all the work that went into Celery, and the other features Celery implements like monitoring, workflows, and a large list of other things that you don't immediately think of when starting a project. It keeps a repository of these patterns for the Python community, and even something like crossbar.io could be supported as a transport.
Celery also supports pub/sub, and other topologies.
>So can your celery workers. It happened to me many times.
With the major difference that your tasks can be redelivered to a different worker, and so will complete anyway.
>Well, the main point of the article is that celery is solving >the GIL. It's not, it's bypassing it
I was agreeing with you there, but I guess my reply was not clear on that. I just wanted to point out some inaccuracies in your reply.
>coroutines and multiprocessing
Be careful using the multiprocessing module, it has some very serious bugs. I've spent the last 4 years rewriting parts of it for Celery
Celery can use eventlet/gevent instead of multiprocessing for executing tasks, so this should be possible (granted, not sure if using it as a web server is a great idea)
>- Tasks cannot communicate with each others;
This is not true, they can send messages to each other
>- You must juggle with the workflow of your tasks (is it ready >? it it dead ?). You can use await stuff() with a try/except;
If you have to juggle it means your workflow design is not good enough. I.e. you shouldn't wait for other tasks, you should have callbacks and errbacks.
Also, note that a try/except does not guarantee the operation will be completed (e.g your asyncio app can be killed).
I'm not sure there's much point in comparing these, they are wildly different concepts: Celery is a distributed system, asyncio is for async I/O, the GIL is only a problem if you can't start n instances of your app.