Per-point ERA5 fetches were tragically slow because every point-hour triggered its own asynchronous CDS job (submit → poll → assemble → download). For backfill this meant thousands of independent jobs queued against Copernicus. The new path groups requests by calendar month and a 2° × 2° lat/lon tile so one CDS cycle populates ~60k profiles at once, and Oban uniqueness on (year, month, tile_lat, tile_lon) collapses every duplicate enqueue. - Era5BatchClient builds the monthly CDS requests, extracts every (lat, lon, hour) from the GRIB2 blob with wgrib2, derives refractivity params, and bulk-inserts in 2k-row chunks with on_conflict: :nothing. fetch_month_into_db/1 short-circuits when the month-tile already has any cached profile. - Era5MonthBatchWorker runs the batch on the :era5 queue with a generous backoff (10m → 1d) and the uniqueness key above. - Era5FetchWorker is now a thin router: cache hit → :ok, cache miss → enqueue the month-batch for the point's tile-month and return. No more per-point CDS calls. - Wgrib2 grows extract_grid_messages/3 which preserves per-message datetimes by parsing `d=YYYYMMDDHH[MMSS]` from the inventory, so a single GRIB2 file carrying a whole month decodes correctly. - The era5_backfill mix task enqueues month-tile batches directly.
298 lines
9.5 KiB
Elixir
298 lines
9.5 KiB
Elixir
defmodule Microwaveprop.Weather.Era5BatchClient do
|
||
@moduledoc """
|
||
Fetches ERA5 reanalysis data from Copernicus Climate Data Store (CDS) one
|
||
month-tile at a time and bulk-inserts the resulting profiles into
|
||
`era5_profiles`.
|
||
|
||
The per-point `Era5Client.fetch_profile/3` path is tragically slow because
|
||
every point-hour triggers its own asynchronous CDS job (submit → poll →
|
||
assemble → download). For backfill this balloons into thousands of
|
||
independent CDS jobs. This module instead groups requests by month and a
|
||
small lat/lon tile (2° × 2° by default) so one submit/poll/download cycle
|
||
produces many thousands of profile rows.
|
||
|
||
Idempotency: `fetch_month_into_db/1` is a no-op when the requested
|
||
(year, month, tile) already has profile rows in `era5_profiles`. Combined
|
||
with Oban uniqueness on `Era5MonthBatchWorker`, redundant fetches are
|
||
collapsed.
|
||
"""
|
||
|
||
import Ecto.Query
|
||
|
||
alias Microwaveprop.Repo
|
||
alias Microwaveprop.Weather.Era5Client
|
||
alias Microwaveprop.Weather.Era5Profile
|
||
alias Microwaveprop.Weather.Grib2.Wgrib2
|
||
alias Microwaveprop.Weather.SoundingParams
|
||
|
||
require Logger
|
||
|
||
# ERA5 hourly reanalysis runs on a 0.25° global grid. Tiles are aligned
|
||
# to an even-degree lattice (2° × 2°), producing 9 × 9 = 81 grid points
|
||
# per tile per hour. A month-tile expands to ~60k profiles.
|
||
@tile_degrees 2
|
||
@era5_step 0.25
|
||
|
||
@doc """
|
||
Returns the integer tile coordinate `{tile_lat, tile_lon}` for a point,
|
||
where `tile_lat` is the south edge and `tile_lon` is the west edge.
|
||
"""
|
||
def tile_for_point(lat, lon) do
|
||
{floor_to_tile(lat), floor_to_tile(lon)}
|
||
end
|
||
|
||
defp floor_to_tile(value) do
|
||
value
|
||
|> Kernel./(@tile_degrees)
|
||
|> Float.floor()
|
||
|> trunc()
|
||
|> Kernel.*(@tile_degrees)
|
||
end
|
||
|
||
@doc false
|
||
def tile_degrees, do: @tile_degrees
|
||
|
||
@doc """
|
||
Fetches one month of ERA5 data covering the given tile and stores every
|
||
resulting hourly profile in `era5_profiles`. Returns
|
||
`{:ok, inserted_count}` or `{:error, reason}`.
|
||
"""
|
||
def fetch_month_into_db(%{year: year, month: month, tile_lat: tile_lat, tile_lon: tile_lon}) do
|
||
if month_tile_cached?(year, month, tile_lat, tile_lon) do
|
||
Logger.info("ERA5Batch: #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon} already cached")
|
||
{:ok, 0}
|
||
else
|
||
do_fetch_month(year, month, tile_lat, tile_lon)
|
||
end
|
||
end
|
||
|
||
defp do_fetch_month(year, month, tile_lat, tile_lon) do
|
||
area = tile_area(tile_lat, tile_lon)
|
||
days = days_in_month(year, month)
|
||
|
||
Logger.info("ERA5Batch: fetching #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon} (#{days} days)")
|
||
|
||
with {:ok, single_grib} <- fetch_single_level_month(year, month, days, area),
|
||
{:ok, pressure_grib} <- fetch_pressure_level_month(year, month, days, area),
|
||
{:ok, profiles} <- build_month_profiles(tile_lat, tile_lon, single_grib, pressure_grib) do
|
||
inserted = bulk_insert_profiles(profiles)
|
||
|
||
Logger.info("ERA5Batch: stored #{inserted} profiles for #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon}")
|
||
|
||
{:ok, inserted}
|
||
end
|
||
end
|
||
|
||
defp month_tile_cached?(year, month, tile_lat, tile_lon) do
|
||
{first, last} = month_bounds(year, month)
|
||
|
||
Era5Profile
|
||
|> where(
|
||
[p],
|
||
p.valid_time >= ^first and p.valid_time < ^last and
|
||
p.lat >= ^(tile_lat * 1.0) and p.lat <= ^(tile_lat * 1.0 + @tile_degrees) and
|
||
p.lon >= ^(tile_lon * 1.0) and p.lon <= ^(tile_lon * 1.0 + @tile_degrees)
|
||
)
|
||
|> Repo.exists?()
|
||
end
|
||
|
||
defp month_bounds(year, month) do
|
||
{:ok, first} = Date.new(year, month, 1)
|
||
{:ok, first_dt} = DateTime.new(first, ~T[00:00:00], "Etc/UTC")
|
||
last_dt = DateTime.add(first_dt, days_in_month(year, month) * 86_400, :second)
|
||
{first_dt, last_dt}
|
||
end
|
||
|
||
defp days_in_month(year, month), do: Calendar.ISO.days_in_month(year, month)
|
||
|
||
defp tile_area(tile_lat, tile_lon) do
|
||
# CDS area is [north, west, south, east]
|
||
Enum.map([tile_lat + @tile_degrees, tile_lon, tile_lat, tile_lon + @tile_degrees], &(&1 * 1.0))
|
||
end
|
||
|
||
defp pad(n), do: String.pad_leading(Integer.to_string(n), 2, "0")
|
||
|
||
# -- CDS requests ----------------------------------------------------------
|
||
|
||
defp fetch_single_level_month(year, month, days, area) do
|
||
request = %{
|
||
"product_type" => ["reanalysis"],
|
||
"variable" => Era5Client.single_level_vars(),
|
||
"year" => [Integer.to_string(year)],
|
||
"month" => [pad(month)],
|
||
"day" => Enum.map(1..days, &pad/1),
|
||
"time" => Enum.map(0..23, &"#{pad(&1)}:00"),
|
||
"area" => area,
|
||
"data_format" => "grib"
|
||
}
|
||
|
||
Era5Client.submit_and_download("reanalysis-era5-single-levels", request)
|
||
end
|
||
|
||
defp fetch_pressure_level_month(year, month, days, area) do
|
||
request = %{
|
||
"product_type" => ["reanalysis"],
|
||
"variable" => Era5Client.pressure_level_vars(),
|
||
"pressure_level" => Era5Client.pressure_levels(),
|
||
"year" => [Integer.to_string(year)],
|
||
"month" => [pad(month)],
|
||
"day" => Enum.map(1..days, &pad/1),
|
||
"time" => Enum.map(0..23, &"#{pad(&1)}:00"),
|
||
"area" => area,
|
||
"data_format" => "grib"
|
||
}
|
||
|
||
Era5Client.submit_and_download("reanalysis-era5-pressure-levels", request)
|
||
end
|
||
|
||
# -- GRIB extraction -------------------------------------------------------
|
||
|
||
defp grid_spec_for_tile(tile_lat, tile_lon) do
|
||
count = round(@tile_degrees / @era5_step) + 1
|
||
|
||
%{
|
||
lon_start: tile_lon * 1.0,
|
||
lon_count: count,
|
||
lon_step: @era5_step,
|
||
lat_start: tile_lat * 1.0,
|
||
lat_count: count,
|
||
lat_step: @era5_step
|
||
}
|
||
end
|
||
|
||
defp build_month_profiles(tile_lat, tile_lon, single_grib, pressure_grib) do
|
||
grid_spec = grid_spec_for_tile(tile_lat, tile_lon)
|
||
|
||
with {:ok, single_msgs} <-
|
||
Wgrib2.extract_grid_messages(single_grib, ":(TMP|DPT|PRES|UGRD|VGRD|PWAT|HPBL):", grid_spec),
|
||
{:ok, pressure_msgs} <- Wgrib2.extract_grid_messages(pressure_grib, ":(TMP|DPT|HGT):", grid_spec) do
|
||
# Index both message sets by (datetime, lat, lon) → var:level map.
|
||
points_single = index_messages(single_msgs)
|
||
points_pressure = index_messages(pressure_msgs)
|
||
|
||
all_keys =
|
||
MapSet.union(
|
||
MapSet.new(Map.keys(points_single)),
|
||
MapSet.new(Map.keys(points_pressure))
|
||
)
|
||
|
||
profiles =
|
||
all_keys
|
||
|> Enum.map(fn {dt, lat, lon} ->
|
||
single_fields = Map.get(points_single, {dt, lat, lon}, %{})
|
||
pressure_fields = Map.get(points_pressure, {dt, lat, lon}, %{})
|
||
build_profile_attrs(lat, lon, dt, single_fields, pressure_fields)
|
||
end)
|
||
|> Enum.reject(&is_nil/1)
|
||
|
||
{:ok, profiles}
|
||
end
|
||
end
|
||
|
||
defp index_messages(messages) do
|
||
Enum.reduce(messages, %{}, fn msg, acc ->
|
||
key_name = "#{msg.var}:#{msg.level}"
|
||
|
||
Enum.reduce(msg.values, acc, fn {{lat, lon}, value}, inner ->
|
||
Map.update(
|
||
inner,
|
||
{msg.datetime, lat, lon},
|
||
%{key_name => value},
|
||
&Map.put(&1, key_name, value)
|
||
)
|
||
end)
|
||
end)
|
||
end
|
||
|
||
defp build_profile_attrs(_lat, _lon, nil, _single, _pressure), do: nil
|
||
|
||
defp build_profile_attrs(lat, lon, datetime, single_fields, pressure_fields) do
|
||
temp_c = kelvin_to_celsius(single_fields["TMP:2 m above ground"])
|
||
dewpoint_c = kelvin_to_celsius(single_fields["DPT:2 m above ground"])
|
||
pressure_mb = pa_to_mb(single_fields["PRES:surface"])
|
||
hpbl_m = single_fields["HPBL:surface"]
|
||
pwat_mm = single_fields["PWAT:entire atmosphere"]
|
||
|
||
profile =
|
||
Era5Client.pressure_levels()
|
||
|> Enum.map(fn level_str ->
|
||
level = String.to_integer(level_str)
|
||
pres_key = "#{level} mb"
|
||
|
||
%{
|
||
"pres" => level * 1.0,
|
||
"tmpc" => kelvin_to_celsius(pressure_fields["TMP:#{pres_key}"]),
|
||
"dwpc" => kelvin_to_celsius(pressure_fields["DPT:#{pres_key}"]),
|
||
"hght" => geopotential_to_height(pressure_fields["HGT:#{pres_key}"])
|
||
}
|
||
end)
|
||
|> Enum.reject(&is_nil(&1["tmpc"]))
|
||
|
||
if profile == [] and is_nil(temp_c) do
|
||
nil
|
||
else
|
||
params = SoundingParams.derive(profile)
|
||
|
||
base = %{
|
||
valid_time: datetime,
|
||
lat: lat,
|
||
lon: lon,
|
||
profile: profile,
|
||
hpbl_m: hpbl_m,
|
||
pwat_mm: pwat_mm,
|
||
surface_temp_c: temp_c,
|
||
surface_dewpoint_c: dewpoint_c,
|
||
surface_pressure_mb: pressure_mb
|
||
}
|
||
|
||
if params do
|
||
Map.merge(base, %{
|
||
surface_refractivity: params.surface_refractivity,
|
||
min_refractivity_gradient: params.min_refractivity_gradient,
|
||
ducting_detected: params.ducting_detected || false,
|
||
duct_characteristics: params.duct_characteristics
|
||
})
|
||
else
|
||
base
|
||
end
|
||
end
|
||
end
|
||
|
||
defp kelvin_to_celsius(nil), do: nil
|
||
defp kelvin_to_celsius(k), do: k - 273.15
|
||
|
||
defp pa_to_mb(nil), do: nil
|
||
defp pa_to_mb(pa), do: pa / 100.0
|
||
|
||
defp geopotential_to_height(nil), do: nil
|
||
defp geopotential_to_height(gp), do: gp / 9.80665
|
||
|
||
# -- Bulk insert -----------------------------------------------------------
|
||
|
||
# Insert in chunks to keep each query under Postgres' parameter limit
|
||
# (65,535 params / ~15 fields ≈ 4,000 rows per statement).
|
||
@chunk_size 2_000
|
||
|
||
defp bulk_insert_profiles(profiles) do
|
||
now = DateTime.truncate(DateTime.utc_now(), :second)
|
||
|
||
profiles
|
||
|> Enum.chunk_every(@chunk_size)
|
||
|> Enum.reduce(0, fn chunk, acc ->
|
||
rows =
|
||
Enum.map(chunk, fn attrs ->
|
||
attrs
|
||
|> Map.put(:inserted_at, now)
|
||
|> Map.put(:updated_at, now)
|
||
end)
|
||
|
||
{count, _} =
|
||
Repo.insert_all(Era5Profile, rows,
|
||
on_conflict: :nothing,
|
||
conflict_target: [:lat, :lon, :valid_time]
|
||
)
|
||
|
||
acc + count
|
||
end)
|
||
end
|
||
end
|