diff --git a/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex b/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex index 925e94ae..40a30b28 100644 --- a/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex +++ b/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex @@ -12,6 +12,8 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do @asos_radius_km 150 @sounding_radius_km 300 + # Oban jobs have ~15 params each; PG limit is 65535 params per query + @insert_batch_size 1000 @impl Oban.Worker def perform(%Oban.Job{}) do @@ -34,7 +36,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do jobs = build_weather_jobs(qsos) if jobs != [] do - Oban.insert_all(jobs) + insert_all_chunked(jobs) end qso_ids = Enum.map(qsos, & &1.id) @@ -53,7 +55,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do jobs = build_hrrr_jobs(qsos) if jobs != [] do - Oban.insert_all(jobs) + insert_all_chunked(jobs) end qso_ids = Enum.map(qsos, & &1.id) @@ -72,7 +74,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do jobs = build_terrain_jobs(qsos) if jobs != [] do - Oban.insert_all(jobs) + insert_all_chunked(jobs) end qso_ids = Enum.map(qsos, & &1.id) @@ -128,7 +130,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do jobs = build_iemre_jobs(qsos) if jobs != [] do - Oban.insert_all(jobs) + insert_all_chunked(jobs) end qso_ids = Enum.map(qsos, & &1.id) @@ -213,4 +215,10 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do }) end end + + defp insert_all_chunked(jobs) do + jobs + |> Enum.chunk_every(@insert_batch_size) + |> Enum.each(&Oban.insert_all/1) + end end