idoshamun

idoshamun

I want to build a system that supports user provided multi steps workflows. Every such step may include interacting with a third party service.
Broadway looks like a perfect fit for the producer side but I didn’t see any complex example that includes multi step flows. Most examples do a local simple transformation and that’s it.
My current thought is to start a GenServer or GenStage on every new message that will handle the execution of the flow.

Showing Posts 1 to 10

zachallaun

zachallaun

Can you please provide a concrete example of a workflow you’re trying to achieve?

I want to build a system that supports user provided multi steps workflows. Every such step may include interacting with a third party service.

When you say “user provided,” do you mean the steps are defined by the user, similar to automation tools like Zapier/Make? Or are the steps predefined, but each step requires some input from a user?

idoshamun

idoshamun OP

Yes, exactly something like Zapier.
More concrete example:
Subscribe to a Google Pub/Sub subscription, dispatch http request to a service with the message as payload, split execution based on the response. If true, format the response and send it as a slack message. Else, drop the message or do something else.
I hope that makes sense, let me know if you need more details

joey_the_snake

joey_the_snake

This can all be done pretty easily with one step. Instead of a simple transformation you could have one step that sends the http request and acts on the response. Not sure if there is a need to split it up.

idoshamun

idoshamun OP

This is another option I’m considering my two concerns:

  • Long running tasks may exceed the acknowledgment deadline of the message queue and trigger a retry. This is why I think it should be better to ack immediately the message and handle it myself. I don’t know if it’s possible without separating.
  • I want to manage the execution state in the database. How long each step took and the overall pipeline. And reflect everything in real-time using LiveView. This is why event sourcing feels like the more intuitive approach.

Sorry for not providing this context earlier. :sweat_smile:

SirWerto

SirWerto

Hello :wave:

Even if it looks beautiful, using processes for handling business logic is not a very good pattern.

You can try to go with something like this

def handle_request(request) do
    case request do
          {:type1, some} -> do_something(request)
          {:type2, some} -> do_anotherthing(request)
    end

end

    defp do_something(request) do
    # code
    send_event_to_db()
    end
end

and then spawn a process for each request

The GenStage library is thought to provide back pressure on your system. With the previous example, if you are spawning too many process and your system can’t handle it you can go with a producer [ConsumerSupervisor](ConsumerSupervisor — gen_stage v1.3.2) architecture to limit the pressure that your system is able to handle just limiting the the amount of process it spawns

idoshamun

idoshamun OP

Would you say that an execution of a workflow is a business logic? To be clear I mean the orchestration of a single workflow. It involves an intermediate state management as well.

SirWerto

SirWerto

I would say that if you have in mind a workflow like this:

receive the request
save state to db
interact with a service
save state to db
interact with another service
save to db
...

you can just use a single process to do everything.

Using GenStage with multiple layers passing state between them would only complicate things.

And to define it, you can use a GenServer spawning more GenServers under a Dynamic Supervisor

joey_the_snake

joey_the_snake

Ah ok I understand what you mean now :). The team I work on has done something similar to this. You can use Broadway to listen to the PubSub events and save the message details in a “jobs” database table and then separately create a GenStage pipeline to poll the table for new jobs and perform work and record all the details of it.

idoshamun

idoshamun OP

One GenStage to pull all the jobs right?

joey_the_snake

joey_the_snake

One GenStage pipeline where the producers poll the table and then pass the information to the consumers who will do the processing. And then you can configure the concurrency for the producers/consumers to however much you need.

You might also be able to use a job processing library like Oban. I never used it myself but I believe the general idea is it stores jobs inside of a Postgres table and schedules them/runs them/updates their status.

Where Next? Top

Trending in Questions Top

Blokh
Hey guys, I’ve got a huge CSV ( around 10 GB ) that needs to be processed hourly Do you guys have any suggestions what is the best prac...
New
RSP87
I’m working on a project that simulates the bumbl example in the programming phoenix book. It acts almost like an email client. We have a...
New
kszambelanczyk
Hello! Could someone please give me a help/sample code, how to delete a file from s3 using waffle/waffle_ecto from Phoenix app. I creat...
New
RemyXRenard
I’m seeing that a list inside a Kino.DataTable will be interpreted as a charlist, even if the Kino.configure() is set to charlists: :as_l...
New
velrest
So my question is quite simple and i have found no conclusive answer on forum, google or AI. Should we use :erlang.float for Integer to ...
New
samoloth
Hi, I’ve just set up an application with ash_authentication. There is only magic link strategy for now, so there is no confirmation add o...
New
FlyingNoodle
If a change or preparation module uses Ash.Changeset.get_argument/2 or Ash.Query.get_argument/2 (or any of the other get_argument functio...
New

Other Trending Topics Top

mudasobwa
I am happy to introduce the very α version of the new programming language compiled to BEAM. Welcome Cure. It has literally three kille...
New
garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
New
marciok
Hi there! We created Gust: A task orchestrator inspired by Airflow. For those who have never heard about Aiflow, it’s a Python-based wor...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge & Solve. They are GUI (Emerge) and State management (S...
New
netoum
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews