Fl4m3Ph03n1x
Ecto multiple streams in 1 transaction
Background
PS: the following situation describes an hypothetical scenario, where I own a company that sells things to customers.
I have an Ecto query that is so big, that my machine cannot handle it. With billions of results returned, there is probably not enough RAM in the world that can handle it.
The solution here (or so my research indicates) is to use streams. Streams were made for potentially infinite sets of results, which would fit my use case.
Problem
So lets imagine that I want to delete All users that bought a given item. Maybe that item was not really legal in their country, and now me, the poor guy in IT, has to fix things so the world doesn’t come down crashing.
Naive way:
item_id = "123asdasd123"
purchase_ids =
Purchases
|> where([p], p.item_id == ^item_id)
|> select([p], p.id)
|> Repo.all()
Users
|> where([u], u.purchase_id in ^purchase_ids)
|> Repo.delete_all()
This is the naive way. I call it naive, because of 2 issues:
- We have so many purchases, that the machine’s memory will overflow (looking at
purchase_idsquery) purchase_idswill likely have more than 100K ids, so the second query (where we delete things) will fail as it hits Postgres parameters limit of 32K: https://stackoverflow.com/a/42251312/1337392
What can I say, our product is highly addictive and very well priced!
Our customers simply cant get enough of it. Don’t know why. Nope. No reason comes to mind. None at all.
With these problems in mind, I cannot help my customers and grow my empire, I mean, little home owned business.
I did find this possible solution:
Stream way:
item_id = "123asdasd123"
purchase_ids =
Purchases
|> where([p], p.item_id == ^item_id)
|> select([p], p.id)
stream = Repo.stream(purchase_ids)
Repo.transacion(fn ->
ids = Enum.to_list(stream)
Users
|> where([u], u.purchase_id in ^ids)
|> Repo.delete_all()
end)
Questions
However, I am not convinced this will work:
- I am using
Enum.to_listand saving everything into a variable, placing everything into memory again. So I am not gaining any advantage by usingRepo.stream. - I still have too many
idsfor myRepo.delete_allto work without blowing up
I guess the one advantage here is that this now a transaction, so either everything goes or nothing goes.
So, the following questions arise:
- How do I properly make use of
streamsin this scenario? - Can I delete items by streaming parameters (
ids) or do I have to manually batch them? - Can I stream ids to
Repo.delete_all?
Marked As Solved
benwilson512
The SQL for this would be:
DELETE FROM users u
USING purchases p
WHERE u.purchase_id = p.id and p.item_id = $1
In Ecto you can do this as simply as:
query = from u in Users,
join: p in assoc(u, :purchase),
where: p.item_id == ^item_id
Repo.delete_all(query)
Also Liked
benwilson512
What you’re missing is that it is only returned if you ask it to be returned. By default, it returns nothing at all except a count of the entities deleted.
LostKobrakai
What @benwilson512 suggested is even better than a subquery, but for more information on subqueries see:
benwilson512
Yeah I’m making some assumptions about your User schema cause it isn’t posted here. I’m assuming that if you have a purchase_id then you’d have a belongs_to :purchase, Purchase. If that isn’t true then you’ll need to change the join.
benwilson512
Notably, you can just write the join out by hand if you have the column and can’t change the schema:
query = from u in Users,
join: p in Purchase, on: p.id == u.purchase_id,
where: p.item_id == ^item_id
LostKobrakai
I know it might not be the solution you’re looking for, but instead of two separate queries you could use subqueries to make them a single, but nested query. No need to send the results of the first query back to elixir in the first place.







