QSO enrichment now groups all path points by HRRR hour and creates one batch job per hour instead of one job per point. The batch job downloads the GRIB2 data once and extracts all needed points from the same binary. Legacy single-point jobs are still supported for backward compatibility.
239 lines
5.8 KiB
Elixir
239 lines
5.8 KiB
Elixir
defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do
|
|
@moduledoc false
|
|
use Oban.Worker, queue: :enqueue, max_attempts: 3
|
|
|
|
alias Microwaveprop.Radio
|
|
alias Microwaveprop.Weather
|
|
alias Microwaveprop.Weather.HrrrClient
|
|
alias Microwaveprop.Workers.HrrrFetchWorker
|
|
alias Microwaveprop.Workers.IemreFetchWorker
|
|
alias Microwaveprop.Workers.TerrainProfileWorker
|
|
alias Microwaveprop.Workers.WeatherFetchWorker
|
|
|
|
@asos_radius_km 150
|
|
@sounding_radius_km 300
|
|
# Oban jobs have ~149 params each; PG limit is 65535 params per query
|
|
@insert_batch_size 400
|
|
|
|
@impl Oban.Worker
|
|
def perform(%Oban.Job{}) do
|
|
enqueue_weather_jobs()
|
|
enqueue_hrrr_jobs()
|
|
enqueue_terrain_jobs()
|
|
enqueue_iemre_jobs()
|
|
|
|
:ok
|
|
end
|
|
|
|
defp enqueue_weather_jobs do
|
|
case Radio.unprocessed_qsos() do
|
|
[] ->
|
|
:ok
|
|
|
|
qsos ->
|
|
Radio.backfill_distances(qsos)
|
|
|
|
jobs = build_weather_jobs(qsos)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
qso_ids = Enum.map(qsos, & &1.id)
|
|
Radio.mark_weather_queued!(qso_ids)
|
|
|
|
enqueue_weather_jobs()
|
|
end
|
|
end
|
|
|
|
defp enqueue_hrrr_jobs do
|
|
case Radio.unprocessed_hrrr_qsos() do
|
|
[] ->
|
|
:ok
|
|
|
|
qsos ->
|
|
jobs = build_hrrr_jobs(qsos)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
qso_ids = Enum.map(qsos, & &1.id)
|
|
Radio.mark_hrrr_queued!(qso_ids)
|
|
|
|
enqueue_hrrr_jobs()
|
|
end
|
|
end
|
|
|
|
defp enqueue_terrain_jobs do
|
|
case Radio.unprocessed_terrain_qsos() do
|
|
[] ->
|
|
:ok
|
|
|
|
qsos ->
|
|
jobs = build_terrain_jobs(qsos)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
qso_ids = Enum.map(qsos, & &1.id)
|
|
Radio.mark_terrain_queued!(qso_ids)
|
|
|
|
# Continue until all QSOs are enqueued
|
|
enqueue_terrain_jobs()
|
|
end
|
|
end
|
|
|
|
def build_terrain_jobs(qsos) do
|
|
qsos
|
|
|> Enum.map(fn qso ->
|
|
TerrainProfileWorker.new(%{"qso_id" => qso.id})
|
|
end)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_weather_jobs(qsos) do
|
|
qsos
|
|
|> Enum.flat_map(&jobs_for_qso/1)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_hrrr_jobs(qsos) do
|
|
# Group all QSO points by HRRR hour, then create one batch job per hour
|
|
# instead of one job per point. This deduplicates GRIB2 downloads.
|
|
qsos
|
|
|> Enum.flat_map(&hrrr_points_for_qso/1)
|
|
|> Enum.group_by(fn {_point, hour} -> hour end, fn {point, _hour} -> point end)
|
|
|> Enum.map(fn {hour, points} ->
|
|
unique_points =
|
|
points
|
|
|> Enum.uniq()
|
|
|> Enum.map(fn {lat, lon} -> %{"lat" => lat, "lon" => lon} end)
|
|
|
|
HrrrFetchWorker.new(%{
|
|
"points" => unique_points,
|
|
"valid_time" => DateTime.to_iso8601(hour)
|
|
})
|
|
end)
|
|
end
|
|
|
|
def build_iemre_jobs(qsos) do
|
|
qsos
|
|
|> Enum.flat_map(&iemre_job_for_qso/1)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
defp jobs_for_qso(qso) do
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
build_asos_jobs(lat, lon, qso.qso_timestamp) ++
|
|
build_raob_jobs(lat, lon, qso.qso_timestamp)
|
|
end)
|
|
end
|
|
|
|
defp enqueue_iemre_jobs do
|
|
case Radio.unprocessed_iemre_qsos() do
|
|
[] ->
|
|
:ok
|
|
|
|
qsos ->
|
|
jobs = build_iemre_jobs(qsos)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
qso_ids = Enum.map(qsos, & &1.id)
|
|
Radio.mark_iemre_queued!(qso_ids)
|
|
|
|
enqueue_iemre_jobs()
|
|
end
|
|
end
|
|
|
|
defp hrrr_points_for_qso(%{pos1: nil}), do: []
|
|
|
|
defp hrrr_points_for_qso(qso) do
|
|
rounded_time = HrrrClient.nearest_hrrr_hour(qso.qso_timestamp)
|
|
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
{rlat, rlon} = Weather.round_to_hrrr_grid(lat, lon)
|
|
|
|
if Weather.has_hrrr_profile?(rlat, rlon, rounded_time) do
|
|
[]
|
|
else
|
|
[{{rlat, rlon}, rounded_time}]
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp build_asos_jobs(lat, lon, timestamp) do
|
|
start_dt = DateTime.add(timestamp, -2 * 3600, :second)
|
|
end_dt = DateTime.add(timestamp, 2 * 3600, :second)
|
|
|
|
lat
|
|
|> Weather.nearby_stations(lon, "asos", @asos_radius_km)
|
|
|> Enum.reject(fn station ->
|
|
Weather.has_surface_observations?(station.id, start_dt, end_dt)
|
|
end)
|
|
|> Enum.map(fn station ->
|
|
WeatherFetchWorker.new(%{
|
|
"fetch_type" => "asos",
|
|
"station_id" => station.id,
|
|
"station_code" => station.station_code,
|
|
"start_dt" => DateTime.to_iso8601(start_dt),
|
|
"end_dt" => DateTime.to_iso8601(end_dt)
|
|
})
|
|
end)
|
|
end
|
|
|
|
defp iemre_job_for_qso(%{pos1: nil}), do: []
|
|
|
|
defp iemre_job_for_qso(qso) do
|
|
date = DateTime.to_date(qso.qso_timestamp)
|
|
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
{rlat, rlon} = Weather.round_to_iemre_grid(lat, lon)
|
|
|
|
if Weather.has_iemre_observation?(rlat, rlon, date) do
|
|
[]
|
|
else
|
|
[
|
|
IemreFetchWorker.new(%{
|
|
"lat" => rlat,
|
|
"lon" => rlon,
|
|
"date" => Date.to_iso8601(date)
|
|
})
|
|
]
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp build_raob_jobs(lat, lon, timestamp) do
|
|
sounding_times = Weather.sounding_times_around(timestamp)
|
|
|
|
stations = Weather.nearby_stations(lat, lon, "sounding", @sounding_radius_km)
|
|
|
|
for station <- stations,
|
|
sounding_time <- sounding_times,
|
|
not Weather.has_sounding?(station.id, sounding_time) do
|
|
WeatherFetchWorker.new(%{
|
|
"fetch_type" => "raob",
|
|
"station_id" => station.id,
|
|
"station_code" => station.station_code,
|
|
"sounding_time" => DateTime.to_iso8601(sounding_time)
|
|
})
|
|
end
|
|
end
|
|
|
|
defp insert_all_chunked(jobs) do
|
|
jobs
|
|
|> Enum.chunk_every(@insert_batch_size)
|
|
|> Enum.each(&Oban.insert_all/1)
|
|
end
|
|
end
|