prop/lib/microwaveprop/weather/grid_cache.ex
Graham McIntire f1846c0a53
perf: reduce per-pod RSS and HRRR chain wall time
Telemetry showed the application-master process holding ~830 MiB of
terms from warm_grid_cache_from_latest_profile — the data lives in
the app master's heap and never GCs because the process is idle.
Running it in a Task.start lets the terms die with the task.

Mark GridCache, MrmsCache, NexradCache, and ScoreCache ETS tables
:compressed. The scored-band-map and HRRR grid data are map-heavy;
compression trims hundreds of MiB at a few percent CPU cost.

Memoise HRRR .idx responses in Microwaveprop.Cache. Published idx
files are immutable for a model run, but the hourly chain re-fetches
the same URL dozens of times across forecast hours. Cuts ~10s per
repeat out of hrrr_fetch_idx.

Force a garbage collect at the end of HrrrFetchWorker.perform to
reclaim the refc binary heap held from GRIB2 ranges before the Oban
producer hands the process its next job.
2026-04-19 14:56:48 -05:00

158 lines
4.8 KiB
Elixir

defmodule Microwaveprop.Weather.GridCache do
@moduledoc """
Node-local ETS cache of derived HRRR grid rows keyed by `valid_time`. Mirrors
`Microwaveprop.Propagation.ScoreCache` but for the `/weather` map.
The Weather map LiveView calls `latest_weather_grid/1` on mount and every
pan/zoom. Each call otherwise hits the 42M-row partitioned `hrrr_profiles`
table, runs per-row `derive_and_clean` transforms, and returns 3-10k rows.
With this cache those calls become in-memory map iterations.
Each cache entry stores `%{{lat, lon} => derived_row}` so per-point lookups
(used by `weather_point_detail/3`) are O(1). Populated by
`Microwaveprop.Weather.warm_grid_cache/1` after the hourly worker upserts
new HRRR data, fanned out across the cluster via the `"weather:cache"`
PubSub topic so every node stays in sync.
"""
use GenServer
alias Phoenix.PubSub
@table :weather_grid_cache
@lock_table :weather_grid_fill_locks
@topic "weather:cache"
@pubsub Microwaveprop.PubSub
@type row :: %{required(:lat) => float(), required(:lon) => float(), optional(atom()) => any()}
@type bounds :: %{optional(String.t()) => float()}
@spec start_link(keyword()) :: GenServer.on_start()
def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__)
@spec fetch(DateTime.t()) :: {:ok, [row()]} | :miss
def fetch(valid_time) do
case :ets.lookup(@table, valid_time) do
[{_, grid}] -> {:ok, grid_to_list(grid)}
[] -> :miss
end
end
@spec fetch_bounds(DateTime.t(), bounds() | nil) :: {:ok, [row()]} | :miss
def fetch_bounds(valid_time, bounds) do
case :ets.lookup(@table, valid_time) do
[{_, grid}] -> {:ok, grid_to_filtered_list(grid, bounds)}
[] -> :miss
end
end
@spec fetch_point(DateTime.t(), float(), float()) :: {:ok, row()} | :miss
def fetch_point(valid_time, lat, lon) do
case :ets.lookup(@table, valid_time) do
[{_, grid}] ->
case Map.get(grid, {lat, lon}) do
nil -> :miss
row -> {:ok, row}
end
[] ->
:miss
end
end
@spec put(DateTime.t(), [row()]) :: :ok
def put(valid_time, rows) do
grid = list_to_grid(rows)
:ets.insert(@table, {valid_time, grid})
:ok
end
@doc "Insert locally AND broadcast to peer nodes via PubSub."
@spec broadcast_put(DateTime.t(), [row()]) :: :ok
def broadcast_put(valid_time, rows) do
PubSub.broadcast(@pubsub, @topic, {:weather_cache_refresh, valid_time, rows})
:ok
end
@spec latest_valid_time() :: DateTime.t() | nil
def latest_valid_time do
match_spec = [{{:"$1", :_}, [], [:"$1"]}]
case :ets.select(@table, match_spec) do
[] -> nil
times -> Enum.max(times, DateTime)
end
end
@spec prune_older_than(DateTime.t()) :: non_neg_integer()
def prune_older_than(cutoff) do
match_spec = [{{:"$1", :_}, [{:<, :"$1", {:const, cutoff}}], [true]}]
:ets.select_delete(@table, match_spec)
end
@spec clear() :: :ok
def clear do
:ets.delete_all_objects(@table)
:ok
end
@doc """
Atomically claim the right to fill the cache for `valid_time`. Returns
`true` if this caller won the claim and should run the fill; `false` if
another caller is already filling. Prevents N concurrent /weather mounts
after a pod restart from each firing the 15-second `load_weather_grid_from_db`
query and starving the Postgres connection pool.
"""
@spec claim_fill(DateTime.t()) :: boolean()
def claim_fill(valid_time) do
:ets.insert_new(@lock_table, {valid_time, :in_progress})
end
@doc "Release a fill lock claimed via `claim_fill/1`."
@spec release_fill(DateTime.t()) :: :ok
def release_fill(valid_time) do
:ets.delete(@lock_table, valid_time)
:ok
end
@spec sync() :: :ok
def sync do
GenServer.call(__MODULE__, :sync)
end
@impl true
def init(_opts) do
:ets.new(@table, [:set, :named_table, :public, :compressed, read_concurrency: true])
:ets.new(@lock_table, [:set, :named_table, :public])
PubSub.subscribe(@pubsub, @topic)
{:ok, %{}}
end
@impl true
def handle_call(:sync, _from, state), do: {:reply, :ok, state}
@impl true
def handle_info({:weather_cache_refresh, valid_time, rows}, state) do
put(valid_time, rows)
{:noreply, state}
end
def handle_info(_msg, state), do: {:noreply, state}
# ---------- Internal ----------
defp list_to_grid(rows) do
Map.new(rows, fn %{lat: lat, lon: lon} = row -> {{lat, lon}, row} end)
end
defp grid_to_list(grid), do: Enum.map(grid, fn {_, row} -> row end)
defp grid_to_filtered_list(grid, nil), do: grid_to_list(grid)
defp grid_to_filtered_list(grid, %{"south" => s, "north" => n, "west" => w, "east" => e}) do
grid
|> Enum.filter(fn {{lat, lon}, _} ->
lat >= s and lat <= n and lon >= w and lon <= e
end)
|> Enum.map(fn {_, row} -> row end)
end
end