Iex.new

Iex.new

Hi,

in our Team at work we have Hackathons and for next year I am thinking about to introduce Elixir. As we have todo with data processing in our daily work I planning to propose a small app that should import data in parallel from a huge csv file (let say for example a list of persons) into a database. And maybe to transform the data in between.
I think I would also like to show the results in the front end using Pheonix.

During my research I came across Flow and Broadway.
As I have not so much experience with Elixir and never used Flow/Broadway, what could be a good fit for my experimentation ?

I find this article on dev.to which seems to close to what I want to do :

Thank you!

Showing Posts 1 to 10

D4no0

D4no0

How big the CSV is? What is the policy on failure and restarts of the server, meaning do you need persistence in processing?

If you don’t have any special requirements, I would just recommend to use Oban. You can just simply stream the file contents into small to medium jobs, then process them concurrently with Oban without having to care about anything else, you will have features like rate-limiting and persistence out of the box.

Iex.new

Iex.new OP

I think something like 2GB.
I do not have any policies at the moment as I am just starting to think about the “workshop” and things that could be nice to show. As we do not have so much tome (2 days) I would like something small but that also give a good overview of what is possible to do.

joey_the_snake

joey_the_snake

Flow would be a good fit for this

dimitarvp

dimitarvp

If you don’t want to store future and current tasks state you can just get away with NimbleCSV and Task.async_stream/3, very easily:

defmodule IngestCSV do
  alias NimbleCSV.RFC4180, as: YourCSV

  def load(path) do
    path
    |> Path.expand()
    |> File.stream!(read_ahead: 524_288)
    |> YourCSV.parse_stream()
    |> Task.async_stream(fn [name, age, address, email] ->
      # do something with a single CSV record.
      # possibly best to also hand off the results to another service?
      # if you want to just process all records and receive a new list
      # then you should just use Stream.map and Enum.to_list at the end.
    end, max_concurrency: 100, on_timeout: :kill_task, ordered: false)
    |> Stream.run()
  end
end

Flow and Broadway are awesome but for a 2GB file the above will serve you just fine. I’ve processed files up until 17GB or so, if memory serves.

Again though, if each record – or a batch of records – takes more time to process then you’ll need to have more persistent workers where Oban will be much more suited.

Iex.new

Iex.new OP

@dimitarvp, thank you!

Yesterday I was able to implement what I wanted to do with the solution you proposed :slight_smile:
In my process I check with the lib file_system if a new file was detected in a specific directory and if yes then it will proceed. Each line of the CSV file is correctly imported into a MySQL database. For that, I have used Ecto.

When trying to set up Ecto, it was proposed to set up my Repo in the application.ex file like :

 children = [
       FileWatcher,
       Store # is my Repo
 ]

config/config.exs:

import Config

config :file_watch_example, :ecto_repos, [Store]
config :file_watch_example, Store,
  database: "store",
  username: "user",
  password: "passwd,
  hostname: "db",
  pool_size: 100

In the Ecto documentation, the namespace is not the same:

 children = [
    MyApp.Repo,
  ]

In my case, I have tried to define MyApp.Store but this did not work for me.

So I have some questions … :slight_smile: :

  1. Must I have used MyApp.Store instead of only Store ?
  2. in my config/config.exs file the 2 lines seem almost the same, it is possible to refactor them into one?:
config :file_watch_example, :ecto_repos, [Store]
config :file_watch_example, Store,

If I understand correctly the first one is to set Store as available Ecto Repo and the second one is for the database configuration. Correct?

This is the code I ended with based on what you wrote:

defmodule IngestCSV do
  alias NimbleCSV.RFC4180, as: YourCSV

  NimbleCSV.define(YourCSV, separator: ";")

  def load(path) do
    path
    |> Path.expand()
    |> File.stream!(read_ahead: 524_288)
    |> YourCSV.parse_stream()
    |> Task.async_stream(fn [name, bar_code, price, currency] ->
      # do something with a single CSV record.
      # possibly best to also hand off the results to another service?
      # if you want to just process all records and receive a new list
      # then you should just use Stream.map and Enum.to_list at the end.
      FileWatchExample.Product.create_product(%{name: name, bar_code: bar_code, price: price, currency: currency})
    end, max_concurrency: 100, on_timeout: :kill_task, ordered: false)
    |> Stream.run()
  end
end
dimitarvp

dimitarvp

Glad you made it work!

I can offer you something that can accelerate the code further: you can put Stream.chunk_every(500) after YourCSV.parse_stream() and then the Task.async_stream will accept a batch of 500 records (not a single record).

And then you can do batch inserts – provided you don’t do Ecto validation that is. If you need validation then inserting one by one is still better.

Also Task.async_stream’s option max_concurrency will depend on your repo pool size in this case. F.ex. if your repo has a pool of only 20 connections then you should change the max_concurrency to 20 (or even 18-19). So have that in mind as well, it’s important and you might see a lot of timeout failures if max_concurrency is too high.

Iex.new

Iex.new OP

thanks for the hint !
I have tried with 100 000 000 rows but only 62 733 738 have been imported.

I got some errors like:

01:14:22.102 [info] MyXQL.Connection (#PID<0.302.0>) disconnected: ** (DBConnection.ConnectionError) client #PID<0.726857.0> exited

01:14:22.096 [error] MyXQL.Connection (#PID<0.227.0>) failed to connect: ** (MyXQL.Error) (1040) Too many connections

I think I will have to tune better the pool (100) and the number of max_concurrency (100).

I will play a bit with it.

What do you think about the questions I asked in my previous answer regarding the namespace and configuration?

dimitarvp

dimitarvp

If the MySQL instance itself is telling you “too many connections” maybe you should look into its own config – maybe it is not configured to accept that many? I mean you can allow Elixir to connect to 100 separate MySQL connections (on the same server) but if MySQL is configured to e.g. allow maximum 80 (just an example) then you’d get an error like that.

So it’s time to check MySQL’s config itself.

Yes, always scope your app’s modules. Never leave top-level modules unless you’re 100% convinced you’ll never get a collision. Which if you say “I am sure” would be famous last words because who knows who will try to integrate with your app / library in the future…

Yes, correct, hence it’s not possible to remove one or the other. The first setting comes in handy if you have several repositories.

Iex.new

Iex.new OP

Thank you @dimitarvp !

dimitarvp

dimitarvp

I’m curious if you fix the connection errors. Please let us know.

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
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
Onor.io
I have what I’ve heard referred to as a “lookup table” in my database. This is a way of assigning codes to common values. One common lo...
New
Trolleger
What approach to take when sending live updates to “random” users Hi! I have a question, I have a little chat app, and when I create a DM...
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
matt-savvy
Anyone here using Honeybadger? My Honeybadger account is being overwhelmed with noise from some bots. Seeing a lot of Bandit.HTTPError...
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

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
mcass19
ExRatatui lets you cook up rich terminal UIs in Elixir, powered by Rust’s ratatui via Rustler NIFs. Build interactive terminal applicatio...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge &amp; 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
wintermeyer
There are three potential reasons for members of this forum to have a look at https://vutuv.de You are tired or annoyed of LinkedIn. Yo...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews