Four changes sized by measured prod telemetry (83m of spans): 1. propagation queue: 2 → 1 slot per pod. Two concurrent forecast-hour steps per pod stacked HRRR grid + native duct grid + scored band map into ~5-6 GiB RSS, OOM-killing every ~15 min. 3-way parallelism cluster-wide still finishes the chain inside the hourly interval. 2. weather queue: 3 → 1 slot per pod. ASOS backfill was 429-thrashing IEM (1,296 retryable jobs; logs were nothing but 429 backoffs). 3. PropagationGridWorker: skip native-level duct fetch on f01..f18. At ~7-11 min/fh and 18 forecast hours, this was the largest single cost per chain. Forecast hours fall back to derived[:min_refractivity_gradient] from the pressure-level profile. f00 still gets full native-level duct analysis. 4. HrrrClient.download_grib_ranges_to_file: parallelize with Task.async_stream (max_concurrency 8). The file-backed variant was sequential, dominating native-duct fetch time on the remaining f00 path. ~20s → ~3s per call.
424 lines
16 KiB
Elixir
424 lines
16 KiB
Elixir
defmodule Microwaveprop.Workers.PropagationGridWorker do
|
||
@moduledoc """
|
||
Oban worker that downloads HRRR data and computes propagation scores
|
||
across the CONUS grid for all bands, one forecast hour at a time.
|
||
|
||
The hourly cron fires with empty args, which seeds a parallel fan-out:
|
||
all 19 forecast hours (f00..f18) are enqueued as independent jobs
|
||
against the `:propagation` queue. Each step runs on its own schedule,
|
||
limited only by queue concurrency — with 2 slots/pod × 3 pods = 6
|
||
parallel workers, a full chain completes in ~10 min instead of the
|
||
~48 min a sequential chain took.
|
||
|
||
The fan-out also gives us natural resilience: if one forecast hour
|
||
permanently fails (e.g. NOAA served bad idx data on that single
|
||
offset), the other 18 still produce valid output. Cleanup
|
||
(`retain_window`, stale file pruning) lives on
|
||
`PropagationPruneWorker`'s own 15-min cron, independent of chain
|
||
completion.
|
||
"""
|
||
|
||
use Oban.Worker,
|
||
queue: :propagation,
|
||
# Highest priority on the shared :propagation queue. MrmsFetchWorker
|
||
# and PropagationPruneWorker also live here; a backlog of those
|
||
# (e.g. after a deploy, or MRMS firing every 2 minutes) would
|
||
# otherwise starve the hourly chain because Oban dispatches
|
||
# priority-then-FIFO. Explicit so the invariant is visible.
|
||
priority: 0,
|
||
# Higher than the default 3 so a few DynamicLifeline rescues
|
||
# (e.g., a rolling deploy that kills a mid-flight chain step)
|
||
# don't exhaust the chain's retry budget and discard the whole
|
||
# run. Legitimate scoring errors still give up after 5 attempts.
|
||
max_attempts: 5,
|
||
# Deduplicate identical jobs across a 1-hour window. Uniqueness
|
||
# is over the full args set, so the seed (`%{}`) collapses with
|
||
# itself — FreshnessMonitor's 5-minute ticks during a long outage
|
||
# no longer stack 24 redundant f00-f18 chains — while chain steps
|
||
# with distinct `run_time` + `forecast_hour` args remain distinct
|
||
# from each other and from the seed. The period is the same as
|
||
# the hourly cron interval so successive hourly cron fires move
|
||
# past the window naturally.
|
||
unique: [period: 3600, states: [:available, :scheduled, :executing, :retryable]]
|
||
|
||
alias Microwaveprop.Commercial
|
||
alias Microwaveprop.Propagation
|
||
alias Microwaveprop.Propagation.BandConfig
|
||
alias Microwaveprop.Propagation.Grid
|
||
alias Microwaveprop.Propagation.ProfilesFile
|
||
alias Microwaveprop.Propagation.ScoreCache
|
||
alias Microwaveprop.Weather
|
||
alias Microwaveprop.Weather.GridCache
|
||
alias Microwaveprop.Weather.HrrrClient
|
||
alias Microwaveprop.Weather.HrrrNativeClient
|
||
alias Microwaveprop.Weather.NexradClient
|
||
|
||
require Logger
|
||
|
||
# Hard ceiling for one forecast hour. A healthy step is ~8-10 min.
|
||
# 20 min gives 2× headroom for a slow HRRR fetch or scoring batch.
|
||
# Oban kills the executing process on timeout, which closes linked
|
||
# ports and cascades SIGKILL to any child wgrib2 subprocess.
|
||
@run_timeout_ms 20 * 60 * 1000
|
||
|
||
@impl Oban.Worker
|
||
def timeout(_job), do: @run_timeout_ms
|
||
|
||
@max_forecast_hour 18
|
||
|
||
@impl Oban.Worker
|
||
def perform(%Oban.Job{args: %{"forecast_hour" => fh, "run_time" => run_time_iso}}) do
|
||
{:ok, run_time, _} = DateTime.from_iso8601(run_time_iso)
|
||
run_chain_step(run_time, fh)
|
||
end
|
||
|
||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||
seed_chain()
|
||
end
|
||
|
||
defp seed_chain do
|
||
# HRRR takes ~45min to publish after the hour. Use 2 hours ago to
|
||
# ensure availability.
|
||
two_hours_ago = DateTime.add(DateTime.utc_now(), -2, :hour)
|
||
run_time = two_hours_ago |> HrrrClient.nearest_hrrr_hour() |> DateTime.truncate(:second)
|
||
|
||
Logger.info("PropagationGrid: seeding chain run_time=#{run_time}, f00-f#{@max_forecast_hour} (parallel)")
|
||
|
||
# Fan out all 19 forecast hours at once. Each step is independent
|
||
# (different HRRR URL, different valid_time file output) so there's
|
||
# no dependency that requires the old sequential chain pattern.
|
||
# With 2 slots/pod × 3 pods = 6 concurrent workers the chain wall
|
||
# time drops from ~48 min to ~10 min. Cleanup (retain_window,
|
||
# prune) runs on PropagationPruneWorker's own 15-min cron.
|
||
jobs =
|
||
for fh <- 0..@max_forecast_hour do
|
||
new(%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => fh})
|
||
end
|
||
|
||
Oban.insert_all(jobs)
|
||
:ok
|
||
end
|
||
|
||
defp run_chain_step(run_time, fh) do
|
||
t_start = System.monotonic_time(:millisecond)
|
||
started_at = DateTime.utc_now()
|
||
points = Grid.conus_points()
|
||
valid_time = DateTime.add(run_time, fh * 3600, :second)
|
||
|
||
Logger.info("PropagationGrid: chain step run_time=#{run_time} fh=#{fh} (#{length(points)} points)")
|
||
|
||
result = process_forecast_hour(points, run_time, fh, valid_time)
|
||
:erlang.garbage_collect()
|
||
|
||
total_ms = System.monotonic_time(:millisecond) - t_start
|
||
record_timing(run_time, fh, valid_time, started_at, total_ms, result)
|
||
|
||
case result do
|
||
:ok ->
|
||
Phoenix.PubSub.broadcast(
|
||
Microwaveprop.PubSub,
|
||
"propagation:updated",
|
||
{:propagation_updated, [valid_time]}
|
||
)
|
||
|
||
Logger.info("PropagationGrid: fh=#{fh} step finished in #{format_duration(total_ms)}")
|
||
:ok
|
||
|
||
other ->
|
||
other
|
||
end
|
||
end
|
||
|
||
# Persist one row to propagation_run_timings. Swallow DB failures so
|
||
# the instrumentation can never brick the chain — if the table is
|
||
# missing or the connection is flaky, we log and carry on.
|
||
defp record_timing(run_time, fh, valid_time, started_at, duration_ms, result) do
|
||
{status, error} =
|
||
case result do
|
||
:ok -> {:ok, nil}
|
||
other -> {:failed, inspect(other)}
|
||
end
|
||
|
||
attrs = %{
|
||
run_time: DateTime.truncate(run_time, :second),
|
||
forecast_hour: fh,
|
||
valid_time: DateTime.truncate(valid_time, :second),
|
||
started_at: started_at,
|
||
finished_at: DateTime.add(started_at, duration_ms, :millisecond),
|
||
duration_ms: duration_ms,
|
||
status: status,
|
||
error: error
|
||
}
|
||
|
||
case Propagation.record_run_timing(attrs) do
|
||
{:ok, _} ->
|
||
:ok
|
||
|
||
{:error, changeset} ->
|
||
Logger.warning("PropagationGrid: timing insert failed: #{inspect(changeset.errors)}")
|
||
end
|
||
rescue
|
||
e -> Logger.warning("PropagationGrid: timing insert raised: #{inspect(e)}")
|
||
end
|
||
|
||
defp process_forecast_hour(points, run_time, forecast_hour, valid_time) do
|
||
label = "f#{String.pad_leading(Integer.to_string(forecast_hour), 2, "0")}"
|
||
|
||
case timed(label, fn ->
|
||
HrrrClient.fetch_grid(points, run_time, forecast_hour: forecast_hour)
|
||
end) do
|
||
{:ok, grid_data} ->
|
||
# HRRR profiles used to be persisted to the hrrr_profiles table
|
||
# here for AsosAdjustmentWorker to re-score from. That was ~12
|
||
# minutes of JSONB inserts per chain (92k rows × 19 forecast
|
||
# hours) with no user-visible benefit — the scores live in the
|
||
# /data/scores files now, and AsosAdjustmentWorker is disabled.
|
||
# Per-contact HRRR enrichment still uses HrrrFetchWorker, which
|
||
# writes its own `is_grid_point: false` rows.
|
||
|
||
# Native-level duct fetch + wgrib2 pass costs ~7-11 min/hour and
|
||
# is only run on f00. Forecast hours fall back to
|
||
# derived[:min_refractivity_gradient] from the pressure-level
|
||
# profile — coarser (~250 m vs ~10-50 m) but good enough for
|
||
# forecast-hour ducting, which is inherently lower-confidence.
|
||
grid_data =
|
||
if forecast_hour == 0 do
|
||
merge_native_duct_data(grid_data, run_time, forecast_hour)
|
||
else
|
||
grid_data
|
||
end
|
||
|
||
:erlang.garbage_collect()
|
||
|
||
# NEXRAD current-hour composite reflectivity catches fast-moving
|
||
# convective cells between HRRR hourly analyses. Only useful for
|
||
# f00 — forecast hours can't see the future radar image.
|
||
grid_data =
|
||
if forecast_hour == 0 do
|
||
merge_nexrad_data(grid_data, valid_time)
|
||
else
|
||
grid_data
|
||
end
|
||
|
||
# Commercial-link inverse sensor — only meaningful for f00 because
|
||
# the measurement is of the current atmospheric state, not a forecast.
|
||
grid_data =
|
||
if forecast_hour == 0 do
|
||
merge_commercial_link_data(grid_data, valid_time)
|
||
else
|
||
grid_data
|
||
end
|
||
|
||
:erlang.garbage_collect()
|
||
|
||
# Persist the fully-enriched grid_data for every forecast hour
|
||
# so (a) /weather can show forecast-hour data after a pod
|
||
# restart and (b) point_detail can rebuild the factor
|
||
# breakdown for a clicked cell at any forecast hour by
|
||
# re-running the scorer against the stored profile.
|
||
persist_profiles(grid_data, valid_time)
|
||
|
||
# Weather map only shows the analysis hour — f01..f18 are
|
||
# forecast data that /weather doesn't render. Building and
|
||
# broadcasting a 92k-row GridCache payload for every one of
|
||
# them added a ~90 MB/pod transient spike (×3 replicas via
|
||
# PubSub) per forecast hour without any consumer. Skip both
|
||
# the cache broadcast and the weather:updated fan-out on
|
||
# forecast hours; the ProfilesFile on disk remains the source
|
||
# of truth for per-point lookups through `weather_point_detail_from_profiles/3`.
|
||
if forecast_hour == 0 do
|
||
rows = Weather.build_grid_cache_rows(grid_data, valid_time)
|
||
GridCache.broadcast_put(valid_time, rows)
|
||
|
||
Phoenix.PubSub.broadcast(
|
||
Microwaveprop.PubSub,
|
||
"weather:updated",
|
||
{:weather_updated, valid_time}
|
||
)
|
||
end
|
||
|
||
scores = compute_scores(grid_data, valid_time, forecast_hour)
|
||
|
||
case Propagation.replace_scores(scores, valid_time) do
|
||
{:ok, count} ->
|
||
Logger.info("PropagationGrid: #{label} → #{count} scores for #{valid_time}")
|
||
warm_cache(valid_time)
|
||
|
||
# Broadcast progress *after* persistence so the map's
|
||
# pipeline chip only advances to "through +Nh" once that
|
||
# hour is actually readable from the scores file. Emitting
|
||
# this before the fetch would push the chip ahead of the
|
||
# map by the full forecast-hour wall time (~10 minutes).
|
||
Phoenix.PubSub.broadcast(
|
||
Microwaveprop.PubSub,
|
||
"propagation:pipeline",
|
||
{:propagation_pipeline_progress, %{forecast_hour: forecast_hour, valid_time: valid_time}}
|
||
)
|
||
|
||
:ok
|
||
|
||
error ->
|
||
Logger.error("PropagationGrid: #{label} replace failed: #{inspect(error)}")
|
||
error
|
||
end
|
||
|
||
error ->
|
||
Logger.warning("PropagationGrid: #{label} fetch failed: #{inspect(error)}")
|
||
error
|
||
end
|
||
end
|
||
|
||
defp persist_profiles(grid_data, valid_time) do
|
||
timed("profiles", fn ->
|
||
try do
|
||
ProfilesFile.write!(valid_time, grid_data)
|
||
rescue
|
||
e ->
|
||
Logger.warning("PropagationGrid: profiles write failed: #{inspect(e)}")
|
||
end
|
||
end)
|
||
end
|
||
|
||
defp warm_cache(valid_time) do
|
||
Enum.each(BandConfig.all_bands(), fn band ->
|
||
Propagation.warm_cache_and_broadcast(band.freq_mhz, valid_time)
|
||
end)
|
||
|
||
ScoreCache.prune_older_than(DateTime.add(DateTime.utc_now(), -2, :hour))
|
||
end
|
||
|
||
defp timed(label, fun) do
|
||
t0 = System.monotonic_time(:millisecond)
|
||
result = fun.()
|
||
elapsed = System.monotonic_time(:millisecond) - t0
|
||
Logger.info("PropagationGrid: #{label} took #{format_duration(elapsed)}")
|
||
result
|
||
end
|
||
|
||
defp format_duration(ms) when ms < 1000, do: "#{ms}ms"
|
||
defp format_duration(ms), do: "#{Float.round(ms / 1000, 1)}s"
|
||
|
||
defp merge_native_duct_data(grid_data, run_time, forecast_hour) do
|
||
hour_dt = HrrrClient.nearest_hrrr_hour(run_time)
|
||
date = DateTime.to_date(hour_dt)
|
||
hour = hour_dt.hour
|
||
grid_spec = Grid.wgrib2_grid_spec()
|
||
|
||
case timed("native", fn ->
|
||
HrrrNativeClient.fetch_native_duct_grid(date, hour, grid_spec, forecast_hour)
|
||
end) do
|
||
{:ok, duct_grid} ->
|
||
Logger.info("PropagationGrid: merged #{map_size(duct_grid)} native duct cells")
|
||
apply_duct_grid(grid_data, duct_grid)
|
||
|
||
{:error, reason} ->
|
||
Logger.warning("PropagationGrid: native duct fetch failed (continuing without): #{inspect(reason)}")
|
||
grid_data
|
||
end
|
||
end
|
||
|
||
defp apply_duct_grid(grid_data, duct_grid) do
|
||
Map.new(grid_data, fn {point, profile} ->
|
||
case Map.get(duct_grid, point) do
|
||
nil -> {point, profile}
|
||
duct -> {point, Map.merge(profile, duct)}
|
||
end
|
||
end)
|
||
end
|
||
|
||
defp merge_commercial_link_data(grid_data, valid_time) do
|
||
# Precompute per-link degradation once (≤10 SQL queries total).
|
||
# Commercial links cluster around DFW so most grid cells see nil —
|
||
# the per-cell path is now a pure haversine check, not a DB query.
|
||
lookup = Commercial.build_link_lookup(valid_time)
|
||
|
||
{merged, boosted} =
|
||
Enum.reduce(grid_data, {%{}, 0}, fn {{lat, lon} = point, profile}, {acc, count} ->
|
||
case Commercial.link_degradation_from_lookup({lat, lon}, lookup) do
|
||
nil ->
|
||
{Map.put(acc, point, profile), count}
|
||
|
||
degradation ->
|
||
{Map.put(acc, point, Map.put(profile, :commercial_link_degradation, degradation)), count + 1}
|
||
end
|
||
end)
|
||
|
||
if boosted > 0 do
|
||
Logger.info("PropagationGrid: commercial-link degradation available for #{boosted} grid cells")
|
||
end
|
||
|
||
merged
|
||
end
|
||
|
||
defp merge_nexrad_data(grid_data, valid_time) do
|
||
points = Map.keys(grid_data)
|
||
|
||
case timed("nexrad", fn -> NexradClient.fetch_frame(valid_time, points) end) do
|
||
{:ok, observations} ->
|
||
apply_nexrad_observations(grid_data, observations)
|
||
|
||
{:error, reason} ->
|
||
Logger.warning("PropagationGrid: NEXRAD fetch failed (continuing without): #{inspect(reason)}")
|
||
grid_data
|
||
end
|
||
end
|
||
|
||
defp apply_nexrad_observations(grid_data, observations) do
|
||
index = Map.new(observations, fn obs -> {{obs.lat, obs.lon}, obs.max_reflectivity_dbz} end)
|
||
non_zero = Enum.count(index, fn {_pt, dbz} -> dbz > 0 end)
|
||
Logger.info("PropagationGrid: NEXRAD merged (#{non_zero} cells with precip)")
|
||
|
||
Map.new(grid_data, fn {point, profile} ->
|
||
case Map.get(index, point) do
|
||
nil -> {point, profile}
|
||
dbz -> {point, Map.put(profile, :nexrad_max_reflectivity_dbz, dbz)}
|
||
end
|
||
end)
|
||
end
|
||
|
||
defp compute_scores(grid_data, valid_time, forecast_hour) do
|
||
# Algorithm is the primary scorer. `factors` is only populated for
|
||
# f00 (the analysis hour) — forecast hours skip the JSONB write so
|
||
# the scoring+upsert phase can land in under a minute instead of
|
||
# ~4-5 minutes. point_detail on forecast hours returns a nil
|
||
# breakdown, which the UI tolerates.
|
||
compute_scores_algorithm(grid_data, valid_time, forecast_hour == 0)
|
||
end
|
||
|
||
@doc false
|
||
# Public for testing. grid_data is a %{{lat, lon} => profile} map.
|
||
def compute_scores_algorithm(grid_data, valid_time, include_factors?) do
|
||
Microwaveprop.Instrument.span(
|
||
[:propagation_grid, :score_band],
|
||
%{point_count: map_size(grid_data)},
|
||
fn ->
|
||
grid_data
|
||
|> Task.async_stream(
|
||
&score_one_point(&1, valid_time, include_factors?),
|
||
max_concurrency: System.schedulers_online() * 2,
|
||
timeout: 30_000
|
||
)
|
||
|> Stream.flat_map(fn
|
||
{:ok, results} -> results
|
||
{:exit, _reason} -> []
|
||
end)
|
||
|> Enum.to_list()
|
||
end
|
||
)
|
||
end
|
||
|
||
defp score_one_point({{lat, lon}, profile}, valid_time, include_factors?) do
|
||
band_scores = Propagation.score_grid_point(profile, valid_time, lat, lon)
|
||
|
||
Enum.map(band_scores, fn r ->
|
||
%{
|
||
lat: lat,
|
||
lon: lon,
|
||
valid_time: valid_time,
|
||
band_mhz: r.band_mhz,
|
||
score: r.score,
|
||
factors: if(include_factors?, do: r.factors)
|
||
}
|
||
end)
|
||
end
|
||
end
|