Improve ergonomics of batch operations on entire database
Use a stream for better memory management, and return IDs for records which haven't been processed correctly.
This commit is contained in:
@@ -3,12 +3,11 @@ defmodule MusicLibrary.Records.Batch do
|
|||||||
|
|
||||||
alias MusicLibrary.Records.Record
|
alias MusicLibrary.Records.Record
|
||||||
alias MusicLibrary.Repo
|
alias MusicLibrary.Repo
|
||||||
|
import Ecto.Query
|
||||||
|
|
||||||
def refresh_musicbrainz_data do
|
def refresh_musicbrainz_data do
|
||||||
Record
|
run_on_all_records(fn record ->
|
||||||
|> Repo.all()
|
import_musicbrainz_data(record)
|
||||||
|> Enum.each(fn r ->
|
|
||||||
import_musicbrainz_data(r)
|
|
||||||
Process.sleep(1000)
|
Process.sleep(1000)
|
||||||
end)
|
end)
|
||||||
end
|
end
|
||||||
@@ -22,15 +21,34 @@ defmodule MusicLibrary.Records.Batch do
|
|||||||
end
|
end
|
||||||
|
|
||||||
def update_release_ids do
|
def update_release_ids do
|
||||||
Record
|
run_on_all_records(&update_release_ids/1)
|
||||||
|> Repo.all()
|
|
||||||
|> Enum.each(&update_release_ids/1)
|
|
||||||
end
|
end
|
||||||
|
|
||||||
def update_release_ids(record) do
|
def update_release_ids(record) do
|
||||||
record
|
record
|
||||||
|> Record.update_release_ids()
|
|> Record.update_release_ids()
|
||||||
|> Repo.update!()
|
|> Repo.update()
|
||||||
|
end
|
||||||
|
|
||||||
|
defp run_on_all_records(fun) do
|
||||||
|
q = from(r in Record)
|
||||||
|
stream = Repo.stream(q)
|
||||||
|
|
||||||
|
Repo.transaction(fn ->
|
||||||
|
Enum.reduce(stream, [], fn record, acc ->
|
||||||
|
case fun.(record) do
|
||||||
|
{:error, reason} ->
|
||||||
|
Logger.error("Failed to run function on record #{record.id} with #{inspect(reason)}")
|
||||||
|
[record.id | acc]
|
||||||
|
|
||||||
|
:ok ->
|
||||||
|
acc
|
||||||
|
|
||||||
|
{:ok, _record} ->
|
||||||
|
acc
|
||||||
|
end
|
||||||
|
end)
|
||||||
|
end)
|
||||||
end
|
end
|
||||||
|
|
||||||
defp musicbrainz do
|
defp musicbrainz do
|
||||||
|
|||||||
Reference in New Issue
Block a user