grufino

grufino

Streaming: VM freezes for large file

Hello!

I am building a pipeline to process, aggregate and transform data from csv files, then write back to another csv file… I load rows from a 19 column csv file, and with some mathematical operations (map reduce style) write back 30 columns in another csv.
And it was going fine until I tried to upload a 25mb file to the application, 250000 rows, and then I decided to stream all the operations instead of eagerly processing… but now that I’m changing function by function with streams, I faced a problem that I don’t understand why, after only 5 fields created, when I try to write to file the program just freezes and stops writing after a few thousand lines.

I’m streaming every single function so it shouldn’t have any locks as far as I understand, and for the first thousands writes it works fine so I wonder what’s happening, in the erlang observer I can only see usage of resources dropped to near 0 and it doesn’t write to the file anymore.

This is my stream function (before I just load from file), and next is my write function:

def process(stream, field_longs_lats, team_settings) do
    main_stream =
      stream
      # Removing once that don't have timestamp
      |> Stream.filter(fn [time | _tl] -> time != "-" end)
      # Filter all duplicated rows by timestamp
      |> Stream.uniq_by(fn [time | _tl] -> time end)
      |> Stream.map(&Transform.apply_row_tranformations/1)

    cumulative_milli =
      main_stream
      |> Stream.map(fn [_time, milli | _tl] -> milli end)
      |> Statistics.cumulative_sum()

    speeds =
      main_stream
      |> Stream.map(fn [_time, _milli, _lat, _long, pace | _tl] ->
        pace
      end)
      |> Stream.map(&Statistics.get_speed/1)

    cals = Motion.calories_per_timestep(cumulative_milli, cumulative_milli)

    long_stream =
      main_stream
      |> Stream.map(fn [_time, _milli, lat | _tl] -> lat end)

    lat_stream =
      main_stream
      |> Stream.map(fn [_time, _milli, _lat, long | _tl] -> long end)

    x_y_tuples =
      RelativeCoordinates.relative_coordinates(long_stream, lat_stream, field_longs_lats)

    x = Stream.map(x_y_tuples, fn {x, _y} -> x end)
    y = Stream.map(x_y_tuples, fn {_x, y} -> y end)

    [x, y, cals, long_stream, lat_stream]
  end

write:

def write_to_file(keyword_list, file_name) do
    file = File.open!(file_name, [:write, :utf8])

    IO.write(file, V4.empty_v4_headers() <> "\n")

    keyword_list
    |> Stream.zip()
    |> Stream.each(&write_tuple_row(&1, file))
    |> Stream.run()

    File.close(file)
  end

@spec write_tuple_row(tuple(), pid()) :: :ok
  def write_tuple_row(tuple, file) do
    IO.inspect("writing #{inspect(tuple)}")

    row_content =
      Tuple.to_list(tuple)
      |> Enum.map_join(",", fn value -> Transformations.to_string(value) end)

    IO.write(file, row_content <> "\n")
  end

Most Liked

ananthakumaran

ananthakumaran

You could look at the current stack trace of the process Process.info(pid, :current_stacktrace), which might provide some clue.

Where Next?

Popular in Questions Top

lessless
I believe there are people here who are dealing with CSV files import on the daily basis, and since Excel is a really popular tool there ...
New
LegitStack
I’m hoping you guys can give me some general advice and perhaps code examples if you’re feeling up to it. I’m very interested in Elixir,...
New
openscript
Hello! Sorry for this astonishing simple question, but I’m really stuck. I try to set up the intellij-elixir plugin, but I don’t know ho...
New
lk-geimfari
What is most correct way to open, read and parse JSON file with poison? For example if we have example.json file in root of some projec...
New
lastday4you
I wanted to check elixir version in phoenix because i found that my elixir is 1.5 but when i use Enum.chunk_by it said the function is un...
New
stefanchrobot
What’s the safe way to decode a JSON string into a struct? I want to avoid calling String.to_atom. Jason.decode can give me a map with st...
New
vac
Hi, I'm quite new in Elixir and I'm trying to format a string to a PEM format. I have the certificate value like MIIDBTCCAe2...... and ...
New
myronmarston
The Elixir Typespec docs show the following syntax for keyword lists in typespecs: # ... | [key: type] # keyword lis...
New
Fl4m3Ph03n1x
Background Let’s assume I have a typical GenServer that receives messages as requests, does some operation in a DB and returns responses....
New
gonzofish
I’m currently trying to understand how to join three tables using Ecto. All the examples I’ve seen use 2, so maybe I’m just missing somet...
New

Other popular topics Top

senggen
Erlang/OTP 25 [erts-13.2.2] [source] [64-bit] [smp:8:8] [ds:8:8:10] [async-threads:1] 15:22:35.803 [error] gen_event {lager_file_backend...
New
Tee
can someone please explain to me how Enum.reduce works with maps
New
Jim
As a follow up to my earlier question: I have the code compiling and running but not getting a successful login from the rest server. ...
New
stefanchrobot
What’s the safe way to decode a JSON string into a struct? I want to avoid calling String.to_atom. Jason.decode can give me a map with st...
New
alice
Hey, Just curious what are the main benefits of Elixir compared to Clojure? When is Elixir more useful than Clojure and vice versa? Th...
New
qwerescape
Is there a way to get the call stack or stack trace at any point in the code? Not from exceptions, but an expression that returns how the...
New
fireproofsocks
Forgive me if this is obvious, but how does one delete a database record WITHOUT selecting it first? https://hexdocs.pm/ecto/Ecto.Repo.h...
New
9mm
I am constructing a JSON object (map) and I need to conditionally set a field. I’m trying to write proper elixir-way code… and I’m at a l...
New
Qqwy
Original source of discussion: This topic on the Pragmatic Programmers' Functional Web Development with Elixir, OTP, and Phoenix forum. ...
New
jay1
Why is it that the mnesia database isn’t the most preferred database for use in Elixir/Phoenix?
New

We're in Beta

About us Mission Statement