lessless

lessless

I have this challenge of importing (downloading and saving) paginate data.

Currently I’m doing it recursively in a workflow step and I think it’s safe because if it will crash (f. ex. due to remote api error).

I designed this code to resume from the last known transaction on the first fetch, so it should also resume from the last saved transaction in a case of a crash

   def insert(account_id) do
    Workflow.new()
    |> Workflow.put_context(%{account_id: account_id})
    |> Workflow.add_cascade(:trigger_sync, &trigger_sync/1, queue: :transactions_sync)
    |> Workflow.add_cascade(:import, &import_txs/1,
      deps: [:trigger_sync],
      queue: :transactions_sync
    )
    |> Oban.insert_all()
  end

  def import_txs(%{account_id: account_id}) do
    account = BankAccount.by_id(account_id)
    import_txs(account.finexer_id, account)
  end

  def import_txs(%{paging: %{next: nil}}, _account) do
    :ok
  end

  def import_txs(source, account) do
    {:ok, txs_resp} = download_transactions(source, account)
    {_, _} = save_transactions(txs_resp, account)
    import_txs(txs_resp, account)
  end

  def download_transactions(finexer_id, account) when is_binary(finexer_id) do
    APIClient.account_transactions(finexer_id, [{:"timestamp.gte", Transaction.account_checkpoint(account_id)}])
  end

  def download_transactions(txs_resp, _account) when is_map(txs_resp) do
    APIClient.account_transactions(txs_resp, :next_page)
  end

  defp save_transactions(txs_resp, account) do
    txs_resp
    |> Map.get(:data)
    |> Transaction.bulk_import(account)
  end

As a matter of curiosity I wonder if workflows design can support replacing downloading paginated data by dynamically adding jobs.

For example import_txs can import first page, check if there is more (txs_resp.page.next != nil) and then schedule next job with %{source: txs_resp.page.next} as an argument.

Showing Posts 1 to 2

sorentwo

sorentwo

Oban Core Team

That design should work. There’s no “finished” state for a workflow, so you can always append another job. Did you give that a try?

lessless

lessless OP

I haven’t yet. I just don’t have enough understanding of Oban to put together a working solution.
I guess the shortest path to that would be to call add_cascade instead of import_txs inside import_txs/2.



def insert(account_id)
  Workflow.new(id: @sync_workflow)
  ...
end

  def import_txs(%{account_id: account_id}) do
    import_txs(account.finexer_id, account_id)
  end

  def import_txs(%{paging: %{next: nil}}, _account_id) do
    :ok
  end

  def import_txs(source, account_id) do
    account = BankAccount.by_id(account_id)
    {:ok, txs_resp} = download_transactions(source, account)
    {_, _} = save_transactions(txs_resp, account)
  
   @sync_wokflow
   |> Workflow.status() # get workflow
   |> Workflow.add_cascade(:import, &import_txs/1,
      deps: [:trigger_sync],
      queue: :transactions_sync
    )
   |> Oban.insert_all
end

So the last function would import the current page and schedule the next job. But how do I pass txs_resp to the next job?

— All posts loaded —

Where Next? Top

Trending in Questions Top

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
nseaSeb
Hello, I know there is an approach for handling lists that allows for optimized traversal, but I can’t recall the specific method (somet...
New
kpanic
Hi everyone, I am toying with the idea of building a “match maker” for giving personal help to people that wants to start coding. I sta...
New
brecabral
Documentation While reading the Scoped Routes section, I noticed that the documentation currently refers to a problem without explainin...
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
asweet-confluent
I recently noticed that Elixir’s Logger defaults its primary log level to :debug when no :logger, :level application configuration is pre...
New
apz
I’m new to elixir and just tried to install the elixirLS extension for VScode(ium) and it is throwing some errors that I would like help ...
New

Other Trending Topics Top

GenericJam
Edit: 2026 May 15 - This post is archived. Mob is alive!! Main docs: mob v0.7.11 — Documentation A bit of explanation for the slightly c...
New
JesseHerrick
Hey, I’m Jesse and I’m the main contributor behind Dexter, a full-featured, lightning-fast Elixir LSP optimized for large codebases. It s...
New
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
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
mhanberg
Hi everyone! The first release candidate for the Expert language server project is now available! We’ve published a press release detai...
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

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews