jeromedoyle

jeromedoyle

I have an app that downloads lots of large files from S3 throughout the day. I’m using ibrowse and streaming the files to disk. This works well, but cpu usage is pretty high when 5+ files are downloading at once. I know elixir isn’t the best choice for computation heavy tasks, but am I running into the same limitation with file downloads in the sense that the constant stream of messages overloads the cpu?

Showing Posts 1 to 8

xlphs

xlphs

What are you using to download? I use gen_tcp directly and can ingest 100MB/s easily on low end PC, although I buffer the IO a lot. It’s definitely something with your download logic.

jeromedoyle

jeromedoyle OP

I’m using ibrowse with async requests. Here’s what the downloader module looks like. It’s called inside Task.async_stream.

defmodule IbrowseDownloader do
  require Logger
  alias Downloader.Utils

  @doc false
  def download(url, destination, filename, request_headers) do
    config = %{
        url: url,
        destination: destination,
        filename: filename,
        request_headers: request_headers,
        id: nil,
        file: nil,
      }
    with {:ibrowse_req_id, id} <- send_request(config),
         {:ok, file} <- handle_response(%{config | id: id}) do
      {:ok, file}
    end
  end

  @doc false
  def send_request(config) do
    config.url
    |> to_charlist()
    |> :ibrowse.send_req(config.request_headers, :get, [], [
      {:stream_to, self()},
      {:connect_timeout, 60_000},
      {:inactivity_timeout, 60_000},
      {:max_sessions, 1000},
      {:max_pipeline_size, 100_000},
    ], 120 * 60 * 1000)
  end

  @doc false
  def handle_response(config) do
    id = config.id
    receive do
      {:ibrowse_async_headers, ^id, '200', headers} ->
        Logger.debug("Received 200 status")
        {:ok, filename} = Utils.filename(headers, config.url, config.filename)
        {:ok, file} = Utils.create_file(config.destination, filename)
        handle_response(%{config | filename: filename, file: file})

      {:ibrowse_async_headers, _, '302', headers} ->
        Logger.debug("Received 302 status")
        [{'Location', redirect_url}] = Enum.filter(headers, fn({name, _value}) -> name == 'Location' end)
        {:ibrowse_req_id, id} = send_request(%{config | url: redirect_url})
        handle_response(%{config | url: redirect_url, id: id})

      {:ibrowse_async_headers, ^id, status, _headers} ->
        Logger.debug("Received #{status} status")
        {:error, status}

      {:ibrowse_async_response, ^id, {:error, :connection_closed}} ->
        Logger.error("Received :connection_closed")
        {:error, :connection_closed}

      {:ibrowse_async_response, ^id, {:error, :req_timedout}} ->
        Logger.error("Received :req_timedout")
        {:error, :req_timedout}

      {:ibrowse_async_response_timeout, ^id} ->
        Logger.error("Received ibrowse_async_response_timeout")
        {:error, :timeout}

      {:ibrowse_async_response, ^id, chunk} ->
        # Logger.debug("Received chunk. size #{length(chunk)}")
        IO.binwrite(config.file, chunk)
        handle_response(config)

      {:ibrowse_async_response_end, ^id} ->
        Logger.debug("Received end")
        File.close(config.file)
        file = Path.join([config.destination, config.filename])
        {:ok, file}
    end
  end

end
axelson

axelson

Scenic Core Team

Is the result of IbrowseDownloader.download/4 the entire file contents? And are you sending those back to the process that originated the task before storing it to disk? If so that’s a lot of unnecessary messaging overhead and you’d be better off doing the entire “task” within a the Task. So download and store to disk all within the same task as one unit of work. If that’s not it, then it would be helpful to see your code that creates the tasks.

jeromedoyle

jeromedoyle OP

The return result is just an {:ok, file} tuple containing the full path to the downloaded file. The entire body of work is being done inside Task.async_stream. This is the code that initiates the downloads.

Task.async_stream(files, fn {url, destination, filename} ->
  IbrowseDownloader.download(url, destination, filename, [@auth_header])
end, timeout: :infinity, max_concurrency: 20)
|> Enum.to_list()
xlphs

xlphs

I’ve never used ibrowse before. Can you log how often you are processing response? I think you need better rate control over how often you are writing to disk, the buffer size depends on hardware but in general try 64KB or larger. ibrowse should have something to let you control stream rate because gen_tcp allows that, like :inet.setopts(state.socket, active: 1)

dimitarvp

dimitarvp

I also haven’t used ibrowse but I’ve used httpotion several times with Task.async_stream and have been able to download 200 files sumultaneously for hours at a time (and store them to an NFS volume, all on a small VPS: 256MB of RAM) and when me and a colleague watched it remotely with :observer the CPU was getting very slightly excited – 7-8% – with rare spikes to 15% (I am guessing garbage collector kicking in).

But I will agree with @xlphs – when in doubt about if the network is causing you problems, always reach for :gen_tcp first. It gives you 99% clear experience and if everything works well in that code then you either keep it and use it, or start making another module that uses a higher-level library and gradually isolate the problem.

jeromedoyle

jeromedoyle OP

I found what was causing my issues. One of the async responses I was getting was {:ibrowse_async_response, id, {:error, :req_timedout}} which matched the {:ibrowse_async_response, ^id, chunk} clause and thus called IO.binwrite with {:error, :req_timedout}. Once I added a clause to handle this and return an error tuple instead of calling IO.binwrite, cpu usage has dropped drastically. Thanks for the help everyone!

dimitarvp

dimitarvp

Nice catch! Good job. Glad you solved it.

— All posts loaded —

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
matt-savvy
Anyone here using Honeybadger? My Honeybadger account is being overwhelmed with noise from some bots. Seeing a lot of Bandit.HTTPError...
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
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