prop/lib/microwaveprop/workers/propagation_grid_worker.ex
Graham McIntire 5ec66df135
Cut propagation_scores write cost with DELETE+COPY, skip-factors, UNLOGGED
The scoring+upsert phase was ~4m40s per forecast hour and dominated
wall time. Three stacked optimizations attack it from different
angles.

replace_scores/2 is a new hot-path writer that does DELETE WHERE
valid_time = $1 followed by a plain insert_all (no ON CONFLICT
resolution). The chain worker rewrites the full (valid_time, all
bands) slice every forecast hour, so conflict detection was pure
waste. AsosAdjustmentWorker still uses upsert_scores because it
only rewrites the subset of cells near a station.

factors is now nullable. Forecast hours f01-f18 pass factors: nil
so the JSONB encode + toast write is skipped entirely — roughly
halves the data volume per run. point_detail/4 coalesces nil to
an empty map so the JS popup renders without a TypeError, and
scorer_diff only pulls the most recent valid_time that still has
factors (the f00 row).

propagation_scores is now UNLOGGED, so inserts bypass WAL entirely.
Durability tradeoff: an unclean shutdown truncates the table, but
PropagationGridWorker rebuilds it from HRRR every 3h so a lost
table is re-populated within one cron cycle.

Also adds docs/plans/2026-04-14-duckdb-scores-storage.md — a
speculative plan for a flat-file / DuckDB rewrite with explicit
trigger conditions for when to pick it up (partitioning deferred
too; revisit only if these three don't solve it).
2026-04-14 13:47:35 -05:00

390 lines
13 KiB
Elixir
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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.
Each `perform/1` processes a single forecast hour (~810 min of wall
time) and enqueues the next hour in the chain. A cron fire with
empty args seeds the chain at f00; subsequent runs carry
`forecast_hour` + `run_time` args. At f18 the chain stops and the
pruner cleans up old scores.
Splitting by forecast hour is a resilience play: a full sweep takes
~3 hours of wall time, longer than the typical pod-restart interval
on this deployment. Under the old "one big perform" design, any
deploy mid-sweep killed the whole run, and max_attempts would
exhaust without recording an error. Per-hour jobs survive deploys
because Lifeline only needs to rescue a single 10-minute step, and
retries re-fetch just that forecast hour.
"""
use Oban.Worker,
queue: :propagation,
max_attempts: 3
alias Microwaveprop.Commercial
alias Microwaveprop.Propagation
alias Microwaveprop.Propagation.BandConfig
alias Microwaveprop.Propagation.Grid
alias Microwaveprop.Propagation.ScoreCache
alias Microwaveprop.Weather
alias Microwaveprop.Weather.GridCache
alias Microwaveprop.Weather.HrrrClient
alias Microwaveprop.Weather.HrrrNativeClient
alias Microwaveprop.Weather.NexradClient
alias Microwaveprop.Weather.SoundingParams
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
@doc """
Enqueue the next chain step after a successful forecast-hour run.
Returns `{:ok, job}` when a new step is enqueued, or `:final` when
`fh` is already `@max_forecast_hour` so the chain has no more work.
Public so the chain entry point and tests can both exercise the
same enqueue path.
"""
@spec enqueue_next_step(DateTime.t(), non_neg_integer()) :: {:ok, Oban.Job.t()} | :final
def enqueue_next_step(_run_time, fh) when fh >= @max_forecast_hour, do: :final
def enqueue_next_step(%DateTime{} = run_time, fh) when fh >= 0 do
{:ok, _job} =
%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => fh + 1}
|> new()
|> Oban.insert()
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}")
{:ok, _job} =
%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => 0}
|> new()
|> Oban.insert()
:ok
end
defp run_chain_step(run_time, fh) do
t_start = System.monotonic_time(:millisecond)
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()
case result do
:ok ->
Phoenix.PubSub.broadcast(
Microwaveprop.PubSub,
"propagation:updated",
{:propagation_updated, [valid_time]}
)
total_ms = System.monotonic_time(:millisecond) - t_start
Logger.info("PropagationGrid: fh=#{fh} step finished in #{format_duration(total_ms)}")
case enqueue_next_step(run_time, fh) do
:final ->
Weather.prune_old_grid_profiles()
Propagation.prune_old_scores()
Logger.info("PropagationGrid: chain complete for run_time=#{run_time}")
{:ok, _job} ->
:ok
end
:ok
other ->
other
end
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")}"
# Emit progress so the map page's pipeline status chip can show
# which forecast hour is currently being scored ("now", "+1h", …).
Phoenix.PubSub.broadcast(
Microwaveprop.PubSub,
"propagation:pipeline",
{:propagation_pipeline_progress, %{forecast_hour: forecast_hour, valid_time: valid_time}}
)
case timed(label, fn ->
HrrrClient.fetch_grid(points, run_time, forecast_hour: forecast_hour)
end) do
{:ok, grid_data} ->
store_hrrr_profiles(grid_data, valid_time)
:erlang.garbage_collect()
# Fetch native duct metrics and merge into grid_data for scoring
grid_data = merge_native_duct_data(grid_data, run_time, forecast_hour)
: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()
# Build weather cache rows from in-memory grid_data — avoids the ~20s
# JSONB round trip that was crashing the worker before f01f18 could run.
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}
)
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)
: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 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 store_hrrr_profiles(grid_data, valid_time) do
count =
grid_data
|> Stream.flat_map(fn {{lat, lon}, profile} ->
build_profile_attrs_for_grid(lat, lon, valid_time, profile)
end)
|> Stream.chunk_every(500)
|> Enum.reduce(0, fn chunk, acc ->
Weather.upsert_hrrr_profiles_batch(chunk)
acc + length(chunk)
end)
Logger.info("PropagationGrid: stored #{count} HRRR profiles")
end
defp build_profile_attrs_for_grid(lat, lon, valid_time, profile) do
temp = profile.surface_temp_c
if temp && temp > -80 && temp < 60 do
params =
if is_list(profile.profile) and length(profile.profile) >= 3,
do: SoundingParams.derive(profile.profile)
attrs = %{
valid_time: valid_time,
lat: lat,
lon: lon,
run_time: profile.run_time,
profile: profile.profile || [],
hpbl_m: profile.hpbl_m,
pwat_mm: profile.pwat_mm,
surface_temp_c: profile.surface_temp_c,
surface_dewpoint_c: profile.surface_dewpoint_c,
surface_pressure_mb: profile.surface_pressure_mb
}
attrs =
if params do
Map.merge(attrs, %{
surface_refractivity: params.surface_refractivity,
min_refractivity_gradient: params.min_refractivity_gradient,
ducting_detected: params.ducting_detected,
duct_characteristics: params.duct_characteristics
})
else
attrs
end
[attrs]
else
[]
end
end
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
# Compute link degradation once per distinct link cluster and cache by
# point. Commercial links only cluster around DFW so most grid points see
# nil — cheap no-op path dominates.
boosted =
Enum.count(grid_data, fn {{lat, lon}, _profile} ->
degradation = Commercial.link_degradation_at({lat, lon}, valid_time)
not is_nil(degradation)
end)
if boosted > 0 do
Logger.info("PropagationGrid: commercial-link degradation available for #{boosted} grid cells")
end
Map.new(grid_data, fn {{lat, lon} = point, profile} ->
case Commercial.link_degradation_at({lat, lon}, valid_time) do
nil -> {point, profile}
degradation -> {point, Map.put(profile, :commercial_link_degradation, degradation)}
end
end)
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
defp compute_scores_algorithm(grid_data, valid_time, include_factors?) do
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)
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