cjk
Hi there,
I have a problem with ecto timeouts (so it seems). In a Elixir application I wrote a module for importing data. It gets a CSV file, goes through that file line by line (with a Stream.map()), looks if this dataset already exists in the database, updates it or inserts it as a new dataset. Pretty basic.
To make that operation restartable I wrap this whole process in a database transaction. The amount of datasets is pretty large and it can take up to an hour.
I know of Ecto timeouts, and thus I start that transaction with timeout: :infinity to avoid timeout problems. It works in development, but in prod I get strange errors:
16:02:34.435 [error] Postgrex.Protocol (#PID<0.423.0>) disconnected: ** (DBConnection.ConnectionError) ssl send: closed
or
** (exit) an exception was raised:
** (DBConnection.ConnectionError) ssl send: closed
(ecto_sql) lib/ecto/adapters/sql.ex:624: Ecto.Adapters.SQL.raise_sql_call_error/1
(ecto_sql) lib/ecto/adapters/sql.ex:557: Ecto.Adapters.SQL.execute/5
(ecto) lib/ecto/repo/queryable.ex:147: Ecto.Repo.Queryable.execute/4
(ecto) lib/ecto/repo/queryable.ex:18: Ecto.Repo.Queryable.all/3
(ecto) lib/ecto/repo/queryable.ex:66: Ecto.Repo.Queryable.one/3
(termitool) lib/termitool/meta/meta.ex:332: Termitool.Meta.get_user_by/1
(termitool) lib/termitool/meta/meta.ex:439: Termitool.Meta.get_by_username/1
(termitool) lib/termitool/meta/meta.ex:465: Termitool.Meta.username_password_auth/2
Basically connection errors at random places in the application. I can avoid this problem by setting timeout: :infinity in the repo configuration, but this seems wrong and dangerous to me.
The code is basically:
Repo.transaction(fn ->
stream
|> Stream.map(fn {row, idx} -> update_or_create_row(row) end)
end,
timeout: :infinity)
I am using a stream because that’s what I get from the CSV parsing library.
What am I doing wrong?
Best regards,
CK
Trending in Questions
Other Trending Topics
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #deployment
- #library
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #channels
- #elixirconf
- #exunit
- #discussion
- #code-sync
- #javascript
- #podcasts
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #full-time-contract
- #podcasts-by-brainlid
- #ecto-query
- #elixir-ls
- #blog-post
- #phoenix_html
- #iex
- #ai
- #graphql
- #genstage
- #elixirconf-us
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #metaprogramming
- #security
- #hex











Showing Posts 1 to 10- Show Best Posts
- Show All Posts (oldest first)
- Show All Posts (newest first)
dimitarvp
This looks strange to me: your last statement in such a processing pipeline should always be using
Enum;Streamjust returns a function.Are you sure any actual work is done by this code in production?
cjk
Sorry - it‘s an Enum.map, not a Stream.map
Edit: ok, I’m at the PC again.
Sorry, I simplified the code a bit, in fact it is more like this:
The reduce counts and throws away the success results and collects the error results with their changesets.
cjk
Nothing, nobody? Am I missing some important information?
NobbZ
Is it the transaction that times out or the individual query?
As far as I can remember DB lessons during university, individual queries during a transaction might take longer than without the transaction due to the doublebookkeeping of indexes or linear scan of things that are only visible in the transaction but not yet commited to the database. Can you check if the problem persists if you increase the timeout of the individual query?
kokolegorille
There are some tools to leverage high import, for example GenStage, which will provide back pressure support.
I would also reach for insert_all, because it’s more efficient to send one query with 10’000 insert than 10’000 single insert.
Enum.chunk_by can be useful to treat smaller pieces.
But all of this does not fit well within a huge transaction.
cjk
The timeouts appear all over the place in the application during the import job, in parts of the application which have nothing to do with the import and are not wrapped in the transaction. But the queries in the transaction seem not to time out, I did not see the connection errors there.
cjk
insert_allis not a viable solution in this case, I either update existing rows or insert new rows, but not all in one table. For example one CSV row can result in 4 rows in 4 different tables.kokolegorille
Instead of inserting new rows, I keep them in a list, and proceed all in one go.
If You have different tables, You still might keep multiple lists (one per table), and proceed with multiple insert_all.
I had similar constraint (huge textfile, multiple tables, over 1’000’000 records) with update or create. I had to be careful to race condition when processing the file concurrently.
After using tasks, poolboy, I finally setup my import pipeline with GenStage, and a couple of insert_all.
cjk
Hm. Race conditions on the data are a non-issue in this case. But I will overhaul my import pipeline, thanks for your input…
That said, I guess I found the underlying cause for this problem. Your hint with concurrently importing the data gave me the idea: the CSV library uses workers to parallelize the CSV reading and parsing. That means that my
Stream.map()executions are parallelized as well.Parallelized execution means: different ecto processes are used. And since the server has traffic and more cores than my workstation this means: more Stream workers, more user connections and thus the pool (I was using the default size of 15) could be exhausted pretty fast.
And indeed, if I reduce the pool size on my workstation to 2 and the timeout to 1 second, I get the same errors all over the place.
To avoid that I now use the
:calleroption on all repo calls, and now it works like a charm on my dev machine. Have still to test it in production, thoughThis also means: my import did not run in the transaction anyways…
cjk
Dammit.
The errors appeared again, random connection drops by Postgrex:
And this time there was not even an import job running, I just removed the
timeout: :infinityfrom the repo configuration.What’s going on?!