# Anyone building a kafka streaming platform for Julia similar to Faust for Python?

**URL:** <https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513>\
**Category:** Community\
**Tags:** question, package\
**Created:** [May 15, 2020, 9:31am UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513 "2020-05-15T09:31:13Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![schlichtanders](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/schlichtanders/32/32145_2.png) [@schlichtanders](https://discourse.julialang.org/u/schlichtanders)\
**Post date:** [May 15, 2020, 9:31am UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/1 "2020-05-15T09:31:13Z")

</div>

Hi all,

I am very excited about the future of Julia’s streaming usecases. With [Transducers.jl](https://github.com/tkf/Transducers.jl) and [OnlineStats.jl](https://github.com/joshday/OnlineStats.jl) there is already a rich foundation established.

What I am unable to find is a streaming library similar to what [Faust](https://github.com/robinhood/faust) does for Python + Kafka. Of course there are other streaming frameworks out there, but Kafka is certainly among the most widely used platforms.

Is anyone already working on implementing a Julia version of Kafa Streams / Flink / [Faust](https://github.com/robinhood/faust)?

If not I would actually be quite motivated to take Faust as an example and port the ideas to Julia 😉  
Any help, pointers, resources, contacts, etc. is highly welcome so that this is going to be a stable and widely used standard package of the julia ecosystem!

---

<div class="post-metadata">

**Author:** ![dfdx](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/dfdx/32/120_2.png) [@dfdx](https://discourse.julialang.org/u/dfdx)\
**Post date:** [May 15, 2020, 1:39pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/2 "2020-05-15T13:39:52Z")

</div>

In case you start working on it, do you have any vision for managing distributed processing? Are you planning to build something atop Julia’s built-in capabilities or use cluster orchestration tools like Kubernetes?

---

<div class="post-metadata">

**Author:** ![schlichtanders](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/schlichtanders/32/32145_2.png) [@schlichtanders](https://discourse.julialang.org/u/schlichtanders)\
**Post date:** [May 15, 2020, 2:41pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/3 "2020-05-15T14:41:10Z")

</div>

Very good question, actually I haven’t thought where exactly the computation is running.  
My idea so far is as simple as reimplementing Faust in Julia.

Faust can be spawn on several workers and recently there also seem to be a kubernetes example [Change history for Faust 1.2 — Faust 1.9.0 documentation](https://faust.readthedocs.io/en/latest/history/changelog-1.2.html?highlight=kubernetes)  
I haven’t used Faust yet, but Faust seems to be the go-to implementation for Python. I just expect porting from python to be much easier than porting from Flink or Kafka Streams 😉

If possible, of course, Julia’s built-in capabilities for cluster orchestration should be supported as well. I would expect that it is actually one of the simpler points to start with.

---

<div class="post-metadata">

**Author:** ![Ratingulate](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ratingulate/32/9242_2.png) [@Ratingulate](https://discourse.julialang.org/u/Ratingulate)\
**Post date:** [May 15, 2020, 2:57pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/4 "2020-05-15T14:57:27Z")

</div>

Can you build off of [dagger.jl](https://github.com/JuliaParallel/Dagger.jl/pulse), so you don’t have to re-implement scheduler, logic etc? It’s actively developed by @jpsamaroo But not sure how much batch processing assumptions are built in.

Here’s a distributed streaming lib built off of dask, dagger’s python analogue. [GitHub - python-streamz/streamz: Real-time stream processing for python](https://github.com/python-streamz/streamz)

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [May 16, 2020, 12:11pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/5 "2020-05-16T12:11:14Z")

</div>

Thanks for the ping! Dagger could one day be an option for implementing streaming/batch processing, however it’s currently missing the ability to dynamically extend the graph at runtime, which is important for continuous data processing (although I have [a PR up that’s working on supporting this](https://github.com/JuliaParallel/Dagger.jl/pull/117)).

Dagger would also need some mechanism to ensure that streaming tasks that expect to be able to send/receive data between each other are all active at the same time (currently, the scheduler only launches as many tasks as there are worker processes), as well as a mechanism to send/receive data asynchronously between concurrently-executing tasks.

All of these are non-trivial to implement flexibly and in a performant manner, so I’ll have to give some thought to the implementation. I’ve noted these items down in my Todo list, and will try to address them in the next 2-3 months.

---

<div class="post-metadata">

**Author:** ![schlichtanders](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/schlichtanders/32/32145_2.png) [@schlichtanders](https://discourse.julialang.org/u/schlichtanders)\
**Post date:** [March 14, 2022, 4:37pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/6 "2022-03-14T16:37:10Z")

</div>

@jpsamaroo I would like to revive the dream of using Julia for streaming purposes.

What do you think is the state of Dagger.jl for streaming processing?  
I see it has a distributed DataFrame as of now. Awesome!

In addition, what is the current development / roadmap plans with:

- supporting unbounded datasets
- supporting stateful streaming
- supporting checkpointing, including checkpointing the streaming state

---

<div class="post-metadata">

**Author:** ![jar1](https://avatars.discourse-cdn.com/v4/letter/j/c0e974/32.png) [@jar1](https://discourse.julialang.org/u/jar1)\
**Post date:** [March 14, 2022, 8:45pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/7 "2022-03-14T20:45:01Z")

</div>

Does Dagger have a mechanism for back pressure? eg [The History of Credit-based Flow Control (Part 1) | by OneFlow | Medium](https://oneflow2020.medium.com/the-history-of-credit-based-flow-control-part-1-342ec6efe23c)

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [March 25, 2022, 4:55pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/8 "2022-03-25T16:55:04Z")

</div>

Hi @schlichtanders !

> [@schlichtanders](#):
>
> What do you think is the state of Dagger.jl for streaming processing?

It’s definitely possible right now, as long as you manage the inter-task channels yourself (manually allocating and passing `RemoteChannel`s to tasks that need to communicate). This isn’t guaranteed to work as-is in the long-term, since we may want to skip scheduling some tasks based on load, memory, or some other metric (which the scheduler has the right to do per Dagger’s semantics). But when that changes, I intend to add a way to force the scheduler to schedule multiple tasks all at the same time so that they can communicate.

> [@schlichtanders](#):
>
> - supporting unbounded datasets
> - supporting stateful streaming
> - supporting checkpointing, including checkpointing the streaming state

We don’t have a built-in API for these right now. Most likely, we would provide a way to communicate to Dagger that a given task can be executed multiple times on a collection which supports `iterate`, and Dagger would automatically transform this into something efficient and maximally parallel.

Stateful execution and checkpointing would be harder, since the way in which states may transition are problem-dependent, and naive transformation by Dagger might produce incorrect execution ordering.

For now, I would implement these features manually (which Dagger’s current APIs should be sufficient for), and then we can see what issues we run into to inform further development.

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [March 25, 2022, 5:00pm UTC](https://discourse.julialang.org/t/anyone-building-a-kafka-streaming-platform-for-julia-similar-to-faust-for-python/39513/9 "2022-03-25T17:00:05Z")

</div>

Dagger does not have a mechanism for back pressure, because we don’t have an explicit API for actors/stateful services. This would probably be best implemented in a separate package.
