From fc51f8afe8ec6c34b99fd8d6caf693ac5ea44cef Mon Sep 17 00:00:00 2001 From: Graham McInitre Date: Thu, 16 Jul 2026 07:33:03 -0500 Subject: [PATCH] perf: concurrent ASOS fetches + bulk upsert in commercial poll worker - Replace sequential Enum.each with Task.async_stream for station I/O - Replace per-row upsert_surface_observation with batch upsert_surface_observations --- lib/microwaveprop/commercial/poll_worker.ex | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/lib/microwaveprop/commercial/poll_worker.ex b/lib/microwaveprop/commercial/poll_worker.ex index 6f991103..dfb43e88 100644 --- a/lib/microwaveprop/commercial/poll_worker.ex +++ b/lib/microwaveprop/commercial/poll_worker.ex @@ -67,7 +67,12 @@ defmodule Microwaveprop.Commercial.PollWorker do now = DateTime.utc_now() start_dt = DateTime.add(now, -3600, :second) - Enum.each(stations, &fetch_station_weather(&1, start_dt, now)) + stations + |> Task.async_stream(&fetch_station_weather(&1, start_dt, now), + timeout: :infinity, + max_concurrency: System.schedulers_online() + ) + |> Stream.run() end defp fetch_station_weather(station_code, start_dt, now) do @@ -90,9 +95,8 @@ defmodule Microwaveprop.Commercial.PollWorker do end defp ingest_asos(station, station_code, rows) do - rows - |> Enum.filter(& &1.observed_at) - |> Enum.each(&Weather.upsert_surface_observation(station, &1)) + filtered = Enum.filter(rows, & &1.observed_at) + Weather.upsert_surface_observations(station, filtered) Logger.info("Fetched #{length(rows)} ASOS observations for #{station_code}") end