Unify batching logic

This commit is contained in:
Claudio Ortolina
2026-02-10 19:29:54 +00:00
parent 15615d27a4
commit 4ef09b47f4
3 changed files with 38 additions and 65 deletions
+4 -33
View File
@@ -3,52 +3,23 @@ defmodule MusicLibrary.Artists.Batch do
alias MusicLibrary.Artists alias MusicLibrary.Artists
alias MusicLibrary.Artists.ArtistInfo alias MusicLibrary.Artists.ArtistInfo
alias MusicLibrary.Repo alias MusicLibrary.Batch
require Logger
def refresh_musicbrainz_data do def refresh_musicbrainz_data do
run_on_all_artist_infos(fn artist_info -> Batch.run_on_all(from(r in ArtistInfo), "artist_info", fn artist_info ->
Artists.refresh_musicbrainz_data_async(artist_info) Artists.refresh_musicbrainz_data_async(artist_info)
end) end)
end end
def refresh_discogs_data do def refresh_discogs_data do
run_on_all_artist_infos(fn artist_info -> Batch.run_on_all(from(r in ArtistInfo), "artist_info", fn artist_info ->
Artists.refresh_discogs_data_async(artist_info) Artists.refresh_discogs_data_async(artist_info)
end) end)
end end
def refresh_wikipedia_data do def refresh_wikipedia_data do
run_on_all_artist_infos(fn artist_info -> Batch.run_on_all(from(r in ArtistInfo), "artist_info", fn artist_info ->
Artists.refresh_wikipedia_data_async(artist_info) Artists.refresh_wikipedia_data_async(artist_info)
end) end)
end end
defp run_on_all_artist_infos(fun) do
q = from(r in ArtistInfo)
stream = Repo.stream(q, max_rows: 50)
Repo.transaction(
fn ->
Enum.reduce(stream, [], fn artist_info, acc ->
case fun.(artist_info) do
{:error, reason} ->
Logger.error(
"Failed to run function on artist_info #{artist_info.id} with #{inspect(reason)}"
)
[artist_info.id | acc]
:ok ->
acc
{:ok, _artist_info} ->
acc
end
end)
end,
timeout: :infinity
)
end
end end
+31
View File
@@ -0,0 +1,31 @@
defmodule MusicLibrary.Batch do
alias MusicLibrary.Repo
require Logger
def run_on_all(queryable, label, fun) do
stream = Repo.stream(queryable, max_rows: 50)
Repo.transaction(
fn ->
Enum.reduce(stream, [], fn record, acc ->
case fun.(record) do
{:error, reason} ->
Logger.error(
"Failed to run function on #{label} #{record.id} with #{inspect(reason)}"
)
[record.id | acc]
:ok ->
acc
{:ok, _result} ->
acc
end
end)
end,
timeout: :infinity
)
end
end
+3 -32
View File
@@ -1,48 +1,19 @@
defmodule MusicLibrary.Records.Batch do defmodule MusicLibrary.Records.Batch do
import Ecto.Query import Ecto.Query
alias MusicLibrary.Batch
alias MusicLibrary.Records alias MusicLibrary.Records
alias MusicLibrary.Records.Record alias MusicLibrary.Records.Record
alias MusicLibrary.Repo
require Logger
def refresh_musicbrainz_data do def refresh_musicbrainz_data do
run_on_all_records(fn record -> Batch.run_on_all(from(r in Record), "record", fn record ->
Records.refresh_musicbrainz_data_async(record) Records.refresh_musicbrainz_data_async(record)
end) end)
end end
def generate_embeddings do def generate_embeddings do
run_on_all_records(fn record -> Batch.run_on_all(from(r in Record), "record", fn record ->
Records.generate_embedding_async(record) Records.generate_embedding_async(record)
end) end)
end end
defp run_on_all_records(fun) do
q = from(r in Record)
stream = Repo.stream(q, max_rows: 50)
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,
timeout: :infinity
)
end
end end