2.2 KiB
2.2 KiB
Workflows
Workflows compose jobs with dependency relationships: sequential, fan-out, and fan-in patterns.
alias Oban.Pro.Workflow
Workflow.new()
|> Workflow.add(:extract, ExtractWorker.new(%{source: "db"}))
|> Workflow.add(:transform, TransformWorker.new(%{}), deps: :extract)
|> Workflow.add(:notify, NotifyWorker.new(%{}), deps: :extract) # fan-out
|> Workflow.add(:load, LoadWorker.new(%{}), deps: [:transform, :notify]) # fan-in
|> Oban.insert_all()
Patterns
- Sequential:
deps: :previous_job - Fan-out: Multiple jobs depend on the same parent
- Fan-in: One job depends on multiple parents with
deps: [:a, :b, :c]
Accessing upstream results
Use recorded: true workers:
def process(job) do
{:ok, data} = Workflow.fetch_recorded(job, :extract)
end
Cascading Functions
Function captures that automatically receive context and upstream results:
Workflow.new()
|> Workflow.put_context(%{source: source})
|> Workflow.add_cascade(:extract, &extract/1)
|> Workflow.add_cascade(:transform, &transform/1, deps: :extract)
|> Workflow.add_cascade(:load, &load/1, deps: :transform)
|> Oban.insert_all()
def extract(%{source: source}), do: %{records: fetch_records(source)}
def transform(%{extract: %{records: records}}), do: %{data: process(records)}
def load(%{transform: %{data: data}}), do: insert_all(data)
Fan-out with {enumerable, function/2}:
Workflow.add_cascade(:accounts, {account_ids, &sync_account/2})
Sub-workflows
Compose from reusable pieces:
extract_flow =
Workflow.new()
|> Workflow.add(:fetch, FetchWorker.new(%{}))
|> Workflow.add(:parse, ParseWorker.new(%{}), deps: :fetch)
Workflow.new()
|> Workflow.add(:setup, SetupWorker.new(%{}))
|> Workflow.add_workflow(:extract, extract_flow, deps: :setup)
|> Workflow.add(:cleanup, CleanupWorker.new(%{}), deps: :extract)
|> Oban.insert_all()
Unique workflows
Ensure only one active workflow per name with unique: true. Duplicates while the first is still
running come back flagged conflict?: true and aren't inserted:
Workflow.new(unique: true, name: "nightly-etl")
|> Workflow.add(:extract, ExtractWorker.new(%{}))
|> Oban.insert_all()