AWS S3 doesn't support multi-range HTTP requests. When given Range: bytes=0-999, 2000-2999 it ignores the header and returns the full file (200 instead of 206). Download each range separately.
210 lines
5.7 KiB
Elixir
210 lines
5.7 KiB
Elixir
defmodule Microwaveprop.Weather.HrrrClient do
|
|
@moduledoc false
|
|
|
|
alias Microwaveprop.Weather.Grib2.Extractor
|
|
|
|
require Logger
|
|
|
|
@hrrr_base "https://noaa-hrrr-bdp-pds.s3.amazonaws.com"
|
|
|
|
@pressure_levels [1000, 975, 950, 925, 900, 850, 800, 700]
|
|
|
|
@surface_messages [
|
|
%{var: "TMP", level: "2 m above ground"},
|
|
%{var: "DPT", level: "2 m above ground"},
|
|
%{var: "PRES", level: "surface"},
|
|
%{var: "HPBL", level: "surface"},
|
|
%{var: "PWAT", level: "entire atmosphere (considered as a single layer)"}
|
|
]
|
|
|
|
# --- Public API ---
|
|
|
|
def fetch_profile(lat, lon, valid_time) do
|
|
hour_dt = nearest_hrrr_hour(valid_time)
|
|
date = DateTime.to_date(hour_dt)
|
|
hour = hour_dt.hour
|
|
|
|
with {:ok, sfc_data} <- fetch_product(date, hour, :surface, lat, lon),
|
|
{:ok, prs_data} <- fetch_product(date, hour, :pressure, lat, lon) do
|
|
merged = Map.merge(sfc_data, prs_data)
|
|
result = build_profile(merged)
|
|
{:ok, Map.put(result, :run_time, hour_dt)}
|
|
end
|
|
end
|
|
|
|
def nearest_hrrr_hour(dt) do
|
|
total_seconds = dt.minute * 60 + dt.second
|
|
rounded_dt = DateTime.add(dt, -total_seconds, :second)
|
|
|
|
if dt.minute >= 30 do
|
|
DateTime.add(rounded_dt, 3600, :second)
|
|
else
|
|
rounded_dt
|
|
end
|
|
end
|
|
|
|
def hrrr_url(date, hour, product) do
|
|
date_str = Calendar.strftime(date, "%Y%m%d")
|
|
hour_str = hour |> Integer.to_string() |> String.pad_leading(2, "0")
|
|
|
|
file =
|
|
case product do
|
|
:surface -> "hrrr.t#{hour_str}z.wrfsfcf00.grib2"
|
|
:pressure -> "hrrr.t#{hour_str}z.wrfprsf00.grib2"
|
|
end
|
|
|
|
"#{@hrrr_base}/hrrr.#{date_str}/conus/#{file}"
|
|
end
|
|
|
|
def parse_idx(text) do
|
|
text
|
|
|> String.split("\n")
|
|
|> Enum.reject(&(&1 == ""))
|
|
|> Enum.map(fn line ->
|
|
parts = String.split(line, ":", parts: 8)
|
|
|
|
%{
|
|
msg: String.to_integer(Enum.at(parts, 0)),
|
|
offset: String.to_integer(Enum.at(parts, 1)),
|
|
var: Enum.at(parts, 3),
|
|
level: Enum.at(parts, 4)
|
|
}
|
|
end)
|
|
end
|
|
|
|
def byte_ranges_for_messages(idx_entries, wanted) do
|
|
Enum.flat_map(wanted, fn w ->
|
|
idx_entries
|
|
|> Enum.with_index()
|
|
|> Enum.flat_map(fn {entry, idx} ->
|
|
if entry.var == w.var && entry.level == w.level do
|
|
end_offset =
|
|
case Enum.at(idx_entries, idx + 1) do
|
|
nil -> entry.offset + 10_000_000
|
|
next -> next.offset - 1
|
|
end
|
|
|
|
[{entry.offset, end_offset}]
|
|
else
|
|
[]
|
|
end
|
|
end)
|
|
end)
|
|
end
|
|
|
|
def build_profile(parsed) do
|
|
sfc_temp_k = parsed["TMP:2 m above ground"]
|
|
sfc_dpt_k = parsed["DPT:2 m above ground"]
|
|
sfc_pres_pa = parsed["PRES:surface"]
|
|
|
|
profile =
|
|
Enum.flat_map(@pressure_levels, fn level ->
|
|
level_str = "#{level} mb"
|
|
tmp = parsed["TMP:#{level_str}"]
|
|
dpt = parsed["DPT:#{level_str}"]
|
|
hgt = parsed["HGT:#{level_str}"]
|
|
|
|
if tmp && dpt && hgt do
|
|
[
|
|
%{
|
|
"pres" => level * 1.0,
|
|
"tmpc" => tmp - 273.15,
|
|
"dwpc" => dpt - 273.15,
|
|
"hght" => hgt
|
|
}
|
|
]
|
|
else
|
|
[]
|
|
end
|
|
end)
|
|
|
|
%{
|
|
surface_temp_c: if(sfc_temp_k, do: sfc_temp_k - 273.15),
|
|
surface_dewpoint_c: if(sfc_dpt_k, do: sfc_dpt_k - 273.15),
|
|
surface_pressure_mb: if(sfc_pres_pa, do: sfc_pres_pa / 100.0),
|
|
hpbl_m: parsed["HPBL:surface"],
|
|
pwat_mm: parsed["PWAT:entire atmosphere (considered as a single layer)"],
|
|
profile: profile
|
|
}
|
|
end
|
|
|
|
# --- Private ---
|
|
|
|
defp fetch_product(date, hour, product, lat, lon) do
|
|
url = hrrr_url(date, hour, product)
|
|
idx_url = url <> ".idx"
|
|
|
|
wanted =
|
|
case product do
|
|
:surface -> @surface_messages
|
|
:pressure -> pressure_messages()
|
|
end
|
|
|
|
Logger.info("HRRR fetching #{product} idx from #{idx_url}")
|
|
|
|
with {:ok, idx_text} <- fetch_idx(idx_url),
|
|
idx_entries = parse_idx(idx_text),
|
|
ranges = byte_ranges_for_messages(idx_entries, wanted),
|
|
_ = Logger.info("HRRR downloading #{length(ranges)} GRIB ranges for #{product}"),
|
|
{:ok, grib_binary} <- download_grib_ranges(url, ranges) do
|
|
Logger.info("HRRR #{product} downloaded #{byte_size(grib_binary)} bytes, extracting")
|
|
Extractor.extract_points(grib_binary, lat, lon)
|
|
end
|
|
end
|
|
|
|
defp pressure_messages do
|
|
for level <- @pressure_levels, var <- ["TMP", "DPT", "HGT"] do
|
|
%{var: var, level: "#{level} mb"}
|
|
end
|
|
end
|
|
|
|
defp fetch_idx(url) do
|
|
case Req.get(url, req_options()) do
|
|
{:ok, %{status: 200, body: body}} ->
|
|
{:ok, body}
|
|
|
|
{:ok, %{status: status}} ->
|
|
{:error, "HRRR idx HTTP #{status}"}
|
|
|
|
{:error, reason} ->
|
|
{:error, reason}
|
|
end
|
|
end
|
|
|
|
defp download_grib_ranges(_url, []), do: {:ok, <<>>}
|
|
|
|
defp download_grib_ranges(url, ranges) do
|
|
Enum.reduce_while(ranges, {:ok, <<>>}, fn {start, stop}, {:ok, acc} ->
|
|
case Req.get(url, [{:headers, [{"Range", "bytes=#{start}-#{stop}"}]} | req_options()]) do
|
|
{:ok, %{status: 206, body: body}} ->
|
|
{:cont, {:ok, acc <> body}}
|
|
|
|
{:ok, %{status: status}} ->
|
|
{:halt, {:error, "HRRR grib HTTP #{status}"}}
|
|
|
|
{:error, reason} ->
|
|
{:halt, {:error, reason}}
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp req_options do
|
|
defaults = [retry: &retry?/2, max_retries: 5, retry_delay: &retry_delay/1]
|
|
overrides = Application.get_env(:microwaveprop, :hrrr_req_options, [])
|
|
Keyword.merge(defaults, overrides)
|
|
end
|
|
|
|
defp retry?(_request, response) do
|
|
case response do
|
|
%Req.Response{status: status} when status in [429, 500, 502, 503, 504] -> true
|
|
%{__exception__: true} -> true
|
|
_ -> false
|
|
end
|
|
end
|
|
|
|
defp retry_delay(n) do
|
|
base = Integer.pow(2, n) * 1_000
|
|
jitter = :rand.uniform(1_000)
|
|
base + jitter
|
|
end
|
|
end
|