mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-05 14:55:08 +00:00
perf(sports): fetch ESPN date chunks concurrently (#596)
* perf(sports): fetch ESPN date chunks concurrently Since ESPN started rejecting `dates=YYYYMMDD-YYYYMMDD` on 2026-09-15, one season request became a chunk per month -- and a month over the 500-event cap becomes a request per day. A cold college-baseball season is about 130 requests, and they went out one at a time. That is slower than the 20s budget `_update_plugins()` shares across every plugin at startup, so scoreboards were logging `update() timed out` on first run and being deferred to the scheduled tick with nothing on the panel. Measured on a Pi 4 against live ESPN, March+April college baseball (63 requests, 3101 events): 11.2s sequential, 1.6s concurrent. Over a whole boot that moved football-scoreboard, ledmatrix-flights and birdnet-go inside the budget -- 13 plugins deferred before, 10 after. Chunks now go out six at a time, in two passes: months and edge days first, then the days of any month that came back capped. Six keeps the shared Session under requests' default pool_maxsize of 10, so no connection is discarded. Merged events still follow `espn_date_chunks` order -- a capped month's days are spliced back into its own slot -- so the payload does not depend on which request won the race. Request order is no longer significant, so the three tests that pinned it compare the chunks as a set and keep asserting the merged event order, which is the part callers actually see. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(sports): drop capped month payloads before fetching their days Review of the concurrent chunk fetch found it raised the worst-case peak memory more than the concurrency explains. The old loop discarded a month that came back at the 500-event cap the moment it saw it; the rewrite kept every capped month alive in `results`/`slots` until all of their day requests had finished. Measured on a Pi 4 fetching 20260201-20260531 college baseball (four capped months, 5462 events), peak RSS growth over the call: sequential (main) 83 MB concurrent, months retained 121 MB (+43) concurrent, one worker 108 MB -- the retention alone was +25 concurrent, months dropped 98-100 MB (+16) docs/LOW_MEMORY_BOARDS.md puts a 1 GB Pi 3B+ at under 200 MB of headroom, where running out makes the board unreachable until a power cycle, so the difference matters. The remaining +16 MB is six responses parsing at once; three workers saved about 6 MB more, within run-to-run noise, so the worker count stays at six. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * docs(sports): state what ESPN_CHUNK_WORKERS was measured to do, not more The comment claimed the sequential fetch made scoreboards blow the 20s startup update() timeout. A boot on this branch still deferred 12 plugins and timed out baseball-scoreboard while its season fetches took 0.74s and 1.12s: the startup budget is spent on other per-plugin work. Say what was measured -- 17.7s sequential, 2.6-3.3s concurrent -- and nothing else. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
+101
-21
@@ -34,7 +34,9 @@ workaround retires itself if ESPN reverts.
|
|||||||
|
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
from datetime import date, timedelta
|
from datetime import date, timedelta
|
||||||
|
from functools import partial
|
||||||
from typing import Any, Dict, List, Optional, Tuple
|
from typing import Any, Dict, List, Optional, Tuple
|
||||||
|
|
||||||
# Above this, ESPN returns a truncated list instead of an error. See module
|
# Above this, ESPN returns a truncated list instead of an error. See module
|
||||||
@@ -44,11 +46,18 @@ ESPN_MAX_LIMIT = 500
|
|||||||
# How long a rejected range keeps later ranges from being tried as ranges.
|
# How long a rejected range keeps later ranges from being tried as ranges.
|
||||||
RANGE_RETRY_SECONDS = 6 * 60 * 60
|
RANGE_RETRY_SECONDS = 6 * 60 * 60
|
||||||
|
|
||||||
|
# How many chunk requests may be in flight at once. Four busy months of
|
||||||
|
# college baseball are ~130 chunks once each is re-asked day by day: 17.7s one
|
||||||
|
# at a time on a Pi 4, 2.6-3.3s six at a time. Kept under requests' default
|
||||||
|
# pool_maxsize of 10 so the shared Session never has to discard connections.
|
||||||
|
ESPN_CHUNK_WORKERS = 6
|
||||||
|
|
||||||
_range_lock = threading.Lock()
|
_range_lock = threading.Lock()
|
||||||
_ranges_rejected_until = 0.0
|
_ranges_rejected_until = 0.0
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"ESPN_MAX_LIMIT",
|
"ESPN_MAX_LIMIT",
|
||||||
|
"ESPN_CHUNK_WORKERS",
|
||||||
"RANGE_RETRY_SECONDS",
|
"RANGE_RETRY_SECONDS",
|
||||||
"clamp_espn_limit",
|
"clamp_espn_limit",
|
||||||
"parse_espn_date_range",
|
"parse_espn_date_range",
|
||||||
@@ -169,6 +178,55 @@ def merge_scoreboard_payloads(payloads: List[Dict[str, Any]]) -> Dict[str, Any]:
|
|||||||
return merged
|
return merged
|
||||||
|
|
||||||
|
|
||||||
|
def _fetch_one_chunk(
|
||||||
|
session, url: str, params: Dict[str, Any], headers, timeout, logger, chunk: str,
|
||||||
|
) -> Optional[Dict[str, Any]]:
|
||||||
|
"""GET a single ``dates=`` chunk, or None when it failed.
|
||||||
|
|
||||||
|
One bad chunk must not sink the rest of the season, so every error is
|
||||||
|
logged and swallowed here rather than raised to the gather below.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
response = session.get(
|
||||||
|
url,
|
||||||
|
params=dict(params, dates=chunk, limit=ESPN_MAX_LIMIT),
|
||||||
|
headers=headers,
|
||||||
|
timeout=timeout,
|
||||||
|
)
|
||||||
|
response.raise_for_status()
|
||||||
|
return response.json()
|
||||||
|
except Exception as exc: # noqa: BLE001 - see docstring
|
||||||
|
if logger:
|
||||||
|
logger.warning("ESPN chunk %s failed, skipping it: %s", chunk, exc)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _fetch_chunks(
|
||||||
|
session, url: str, params: Dict[str, Any], headers, timeout, logger,
|
||||||
|
chunks: List[str],
|
||||||
|
) -> List[Optional[Dict[str, Any]]]:
|
||||||
|
"""Fetch every chunk, returning payloads positionally aligned with ``chunks``.
|
||||||
|
|
||||||
|
Requests go out ``ESPN_CHUNK_WORKERS`` at a time because a cold season is
|
||||||
|
over a hundred of them. The order they come back in is not significant --
|
||||||
|
callers keep ``chunks`` order from the returned list -- but it does mean
|
||||||
|
the session is shared across threads, which is why this only ever issues
|
||||||
|
GETs and never touches session state.
|
||||||
|
"""
|
||||||
|
if not chunks:
|
||||||
|
return []
|
||||||
|
fetch = partial(
|
||||||
|
_fetch_one_chunk, session, url, params, headers, timeout, logger,
|
||||||
|
)
|
||||||
|
if len(chunks) == 1:
|
||||||
|
return [fetch(chunks[0])]
|
||||||
|
workers = min(ESPN_CHUNK_WORKERS, len(chunks))
|
||||||
|
with ThreadPoolExecutor(
|
||||||
|
max_workers=workers, thread_name_prefix="espn-chunk",
|
||||||
|
) as pool:
|
||||||
|
return list(pool.map(fetch, chunks))
|
||||||
|
|
||||||
|
|
||||||
def fetch_espn_date_chunks(
|
def fetch_espn_date_chunks(
|
||||||
session,
|
session,
|
||||||
url: str,
|
url: str,
|
||||||
@@ -188,6 +246,11 @@ def fetch_espn_date_chunks(
|
|||||||
largest value that does not corrupt the answer. A month that comes back
|
largest value that does not corrupt the answer. A month that comes back
|
||||||
with 500 events is assumed truncated and re-asked day by day. A failed
|
with 500 events is assumed truncated and re-asked day by day. A failed
|
||||||
chunk is logged and skipped so one bad day cannot cost a whole season.
|
chunk is logged and skipped so one bad day cannot cost a whole season.
|
||||||
|
|
||||||
|
Chunks go out ``ESPN_CHUNK_WORKERS`` at a time, in two passes: the months
|
||||||
|
and edge days first, then the days of any month that came back capped.
|
||||||
|
Merged events keep ``espn_date_chunks`` order regardless of which request
|
||||||
|
finished first, so the result does not depend on the race.
|
||||||
"""
|
"""
|
||||||
params = dict(params or {})
|
params = dict(params or {})
|
||||||
span = parse_espn_date_range(params.get("dates"))
|
span = parse_espn_date_range(params.get("dates"))
|
||||||
@@ -201,35 +264,52 @@ def fetch_espn_date_chunks(
|
|||||||
params.get("dates"), len(chunks),
|
params.get("dates"), len(chunks),
|
||||||
)
|
)
|
||||||
|
|
||||||
payloads: List[Dict[str, Any]] = []
|
results = _fetch_chunks(
|
||||||
attempted = 0
|
session, url, params, headers, timeout, logger, chunks,
|
||||||
pending = list(chunks)
|
)
|
||||||
while pending:
|
attempted = len(chunks)
|
||||||
chunk = pending.pop(0)
|
|
||||||
attempted += 1
|
# A month that came back at the cap is truncated; its days replace it in
|
||||||
try:
|
# place, so merged events stay in chunk order however the requests raced.
|
||||||
response = session.get(
|
slots: List[Any] = results
|
||||||
url,
|
capped: Dict[int, List[str]] = {}
|
||||||
params=dict(params, dates=chunk, limit=ESPN_MAX_LIMIT),
|
for index, chunk in enumerate(chunks):
|
||||||
headers=headers,
|
payload = slots[index]
|
||||||
timeout=timeout,
|
if payload is None or len(chunk) != 6:
|
||||||
)
|
|
||||||
response.raise_for_status()
|
|
||||||
payload = response.json()
|
|
||||||
except Exception as exc: # noqa: BLE001 - one bad chunk must not sink the rest
|
|
||||||
if logger:
|
|
||||||
logger.warning("ESPN chunk %s failed, skipping it: %s", chunk, exc)
|
|
||||||
continue
|
continue
|
||||||
events = payload.get("events") if isinstance(payload, dict) else None
|
events = payload.get("events") if isinstance(payload, dict) else None
|
||||||
if len(chunk) == 6 and len(events or []) >= ESPN_MAX_LIMIT:
|
if len(events or []) >= ESPN_MAX_LIMIT:
|
||||||
if logger:
|
if logger:
|
||||||
logger.info(
|
logger.info(
|
||||||
"ESPN month %s hit the %d-event cap; re-asking it day by day",
|
"ESPN month %s hit the %d-event cap; re-asking it day by day",
|
||||||
chunk, ESPN_MAX_LIMIT,
|
chunk, ESPN_MAX_LIMIT,
|
||||||
)
|
)
|
||||||
pending[:0] = _days_of_month(chunk)
|
capped[index] = _days_of_month(chunk)
|
||||||
|
# Drop the truncated month now rather than after its days arrive:
|
||||||
|
# a capped college-baseball month is ~2MB of parsed JSON, and
|
||||||
|
# holding four of them through ~120 day requests added ~25MB to
|
||||||
|
# the peak -- more than the concurrency itself. Low-memory boards
|
||||||
|
# (docs/LOW_MEMORY_BOARDS.md) have under 200MB of headroom.
|
||||||
|
slots[index] = None
|
||||||
|
payload = events = None
|
||||||
|
|
||||||
|
if capped:
|
||||||
|
days = [day for index in sorted(capped) for day in capped[index]]
|
||||||
|
attempted += len(days)
|
||||||
|
by_day = dict(zip(days, _fetch_chunks(
|
||||||
|
session, url, params, headers, timeout, logger, days,
|
||||||
|
)))
|
||||||
|
for index, month_days in capped.items():
|
||||||
|
slots[index] = [by_day.get(day) for day in month_days]
|
||||||
|
|
||||||
|
payloads: List[Dict[str, Any]] = []
|
||||||
|
for slot in slots:
|
||||||
|
if slot is None:
|
||||||
continue
|
continue
|
||||||
payloads.append(payload)
|
if isinstance(slot, list):
|
||||||
|
payloads.extend(payload for payload in slot if payload is not None)
|
||||||
|
else:
|
||||||
|
payloads.append(slot)
|
||||||
|
|
||||||
if not payloads:
|
if not payloads:
|
||||||
return None
|
return None
|
||||||
|
|||||||
@@ -115,8 +115,11 @@ def test_a_rejected_season_is_recovered_and_cached(service, cache):
|
|||||||
def test_a_full_season_costs_chunks_not_one_request_per_day(service):
|
def test_a_full_season_costs_chunks_not_one_request_per_day(service):
|
||||||
session = RangeRejectingSession({"202609": [{"id": "a"}]})
|
session = RangeRejectingSession({"202609": [{"id": "a"}]})
|
||||||
submit_and_wait(service, session, "20260801-20270301")
|
submit_and_wait(service, session, "20260801-20270301")
|
||||||
assert [call["dates"] for call in session.calls] == [
|
sent = [call["dates"] for call in session.calls]
|
||||||
"20260801-20270301",
|
# Eight chunks rather than 213 per-day requests. They are fetched
|
||||||
|
# concurrently, so the range is the only one pinned to a position.
|
||||||
|
assert sent[0] == "20260801-20270301"
|
||||||
|
assert sorted(sent[1:]) == [
|
||||||
"202608",
|
"202608",
|
||||||
"202609",
|
"202609",
|
||||||
"202610",
|
"202610",
|
||||||
|
|||||||
+123
-4
@@ -12,6 +12,8 @@ Nothing here touches the network. The fake session records what a caller would
|
|||||||
have sent, which is the part that regressed.
|
have sent, which is the part that regressed.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
from datetime import date, timedelta
|
from datetime import date, timedelta
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
@@ -245,7 +247,10 @@ class TestFetch:
|
|||||||
)
|
)
|
||||||
assert [e["id"] for e in data["events"]] == ["a", "b", "c"]
|
assert [e["id"] for e in data["events"]] == ["a", "b", "c"]
|
||||||
sent = [call["dates"] for call in session.calls]
|
sent = [call["dates"] for call in session.calls]
|
||||||
assert sent == ["20260901-20261001", "202609", "20261001"]
|
# Chunks race, so only the rejected range is pinned to a position --
|
||||||
|
# the merged event order above is what has to stay deterministic.
|
||||||
|
assert sent[0] == "20260901-20261001"
|
||||||
|
assert sorted(sent[1:]) == ["202609", "20261001"]
|
||||||
|
|
||||||
def test_chunk_requests_keep_the_clamped_limit(self):
|
def test_chunk_requests_keep_the_clamped_limit(self):
|
||||||
session = FakeSession({"202609": []})
|
session = FakeSession({"202609": []})
|
||||||
@@ -301,9 +306,11 @@ class TestMonthCap:
|
|||||||
|
|
||||||
data = fetch_espn_date_chunks(session, URL, params={"dates": "20260301-20260331"})
|
data = fetch_espn_date_chunks(session, URL, params={"dates": "20260301-20260331"})
|
||||||
|
|
||||||
assert [call["dates"] for call in session.calls] == ["202603"] + [
|
sent = [call["dates"] for call in session.calls]
|
||||||
"202603%02d" % day for day in range(1, 32)
|
# The month has to be asked before its days can be known to be needed;
|
||||||
]
|
# the days themselves race, so compare them as a set.
|
||||||
|
assert sent[0] == "202603"
|
||||||
|
assert sorted(sent[1:]) == ["202603%02d" % day for day in range(1, 32)]
|
||||||
# The truncated month payload is dropped, not merged with the days.
|
# The truncated month payload is dropped, not merged with the days.
|
||||||
assert [event["id"] for event in data["events"]] == [
|
assert [event["id"] for event in data["events"]] == [
|
||||||
"d%d" % day for day in range(1, 32)
|
"d%d" % day for day in range(1, 32)
|
||||||
@@ -370,3 +377,115 @@ class TestRejectedRangeMemo:
|
|||||||
with pytest.raises(RuntimeError):
|
with pytest.raises(RuntimeError):
|
||||||
fetch_espn_scoreboard(session, URL, params={"dates": "20260914-20260915"})
|
fetch_espn_scoreboard(session, URL, params={"dates": "20260914-20260915"})
|
||||||
assert session.calls[-1]["dates"] == "20260914-20260915"
|
assert session.calls[-1]["dates"] == "20260914-20260915"
|
||||||
|
|
||||||
|
|
||||||
|
class TestConcurrency:
|
||||||
|
"""Chunks go out in parallel, which must not change what comes back.
|
||||||
|
|
||||||
|
A cold college-baseball season is ~130 chunks once February through May
|
||||||
|
are re-asked day by day. Sequentially that outran the 20s plugin update()
|
||||||
|
timeout on a Pi, so the requests now overlap -- but the merged payload has
|
||||||
|
to stay exactly what the sequential version produced.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def test_events_keep_chunk_order_however_the_requests_race(self):
|
||||||
|
# Answer the later chunks fastest, so completion order is the reverse
|
||||||
|
# of chunk order and a naive gather would interleave them wrongly.
|
||||||
|
class RacingSession(FakeSession):
|
||||||
|
def get(self, url, params=None, headers=None, timeout=None):
|
||||||
|
dates = str((params or {}).get("dates", ""))
|
||||||
|
if len(dates) == 6:
|
||||||
|
time.sleep(0.02 / (int(dates[4:]) or 1))
|
||||||
|
return super().get(url, params=params, headers=headers, timeout=timeout)
|
||||||
|
|
||||||
|
session = RacingSession(
|
||||||
|
{
|
||||||
|
"202609": [{"id": "sep"}],
|
||||||
|
"202610": [{"id": "oct"}],
|
||||||
|
"202611": [{"id": "nov"}],
|
||||||
|
}
|
||||||
|
)
|
||||||
|
data = fetch_espn_date_chunks(
|
||||||
|
session, URL, params={"dates": "20260901-20261130"}
|
||||||
|
)
|
||||||
|
assert [event["id"] for event in data["events"]] == ["sep", "oct", "nov"]
|
||||||
|
|
||||||
|
def test_a_capped_month_splices_its_days_in_place(self):
|
||||||
|
# October is capped and expands to 31 days; September and November
|
||||||
|
# must still bracket those days in the merged result.
|
||||||
|
full = [{"id": "cap%d" % i} for i in range(ESPN_MAX_LIMIT)]
|
||||||
|
by_chunk = {
|
||||||
|
"202609": [{"id": "sep"}],
|
||||||
|
"202610": full,
|
||||||
|
"202611": [{"id": "nov"}],
|
||||||
|
}
|
||||||
|
by_chunk.update(
|
||||||
|
{"202610%02d" % day: [{"id": "oct%02d" % day}] for day in range(1, 32)}
|
||||||
|
)
|
||||||
|
session = FakeSession(by_chunk)
|
||||||
|
|
||||||
|
data = fetch_espn_date_chunks(
|
||||||
|
session, URL, params={"dates": "20260901-20261130"}
|
||||||
|
)
|
||||||
|
|
||||||
|
expected = ["sep"] + ["oct%02d" % day for day in range(1, 32)] + ["nov"]
|
||||||
|
assert [event["id"] for event in data["events"]] == expected
|
||||||
|
|
||||||
|
def test_two_capped_months_expand_without_crossing_over(self):
|
||||||
|
full = [{"id": "cap%d" % i} for i in range(ESPN_MAX_LIMIT)]
|
||||||
|
by_chunk = {"202609": full, "202610": full}
|
||||||
|
by_chunk.update(
|
||||||
|
{"202609%02d" % day: [{"id": "s%02d" % day}] for day in range(1, 31)}
|
||||||
|
)
|
||||||
|
by_chunk.update(
|
||||||
|
{"202610%02d" % day: [{"id": "o%02d" % day}] for day in range(1, 32)}
|
||||||
|
)
|
||||||
|
session = FakeSession(by_chunk)
|
||||||
|
|
||||||
|
data = fetch_espn_date_chunks(
|
||||||
|
session, URL, params={"dates": "20260901-20261031"}
|
||||||
|
)
|
||||||
|
|
||||||
|
expected = ["s%02d" % day for day in range(1, 31)] + [
|
||||||
|
"o%02d" % day for day in range(1, 32)
|
||||||
|
]
|
||||||
|
assert [event["id"] for event in data["events"]] == expected
|
||||||
|
|
||||||
|
def test_a_failed_day_inside_a_capped_month_only_costs_that_day(self):
|
||||||
|
full = [{"id": "cap%d" % i} for i in range(ESPN_MAX_LIMIT)]
|
||||||
|
by_chunk = {"202610": full}
|
||||||
|
by_chunk.update(
|
||||||
|
{"202610%02d" % day: [{"id": "o%02d" % day}] for day in range(1, 32)}
|
||||||
|
)
|
||||||
|
session = FakeSession(by_chunk, fail_chunks={"20261015"})
|
||||||
|
|
||||||
|
data = fetch_espn_date_chunks(
|
||||||
|
session, URL, params={"dates": "20261001-20261031"}
|
||||||
|
)
|
||||||
|
|
||||||
|
expected = ["o%02d" % day for day in range(1, 32) if day != 15]
|
||||||
|
assert [event["id"] for event in data["events"]] == expected
|
||||||
|
|
||||||
|
def test_no_more_than_the_worker_cap_are_in_flight_at_once(self):
|
||||||
|
live = {"now": 0, "peak": 0}
|
||||||
|
guard = threading.Lock()
|
||||||
|
|
||||||
|
class CountingSession(FakeSession):
|
||||||
|
def get(self, url, params=None, headers=None, timeout=None):
|
||||||
|
with guard:
|
||||||
|
live["now"] += 1
|
||||||
|
live["peak"] = max(live["peak"], live["now"])
|
||||||
|
try:
|
||||||
|
time.sleep(0.01)
|
||||||
|
return super().get(
|
||||||
|
url, params=params, headers=headers, timeout=timeout
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
with guard:
|
||||||
|
live["now"] -= 1
|
||||||
|
|
||||||
|
session = CountingSession()
|
||||||
|
fetch_espn_date_chunks(session, URL, params={"dates": "20260101-20261231"})
|
||||||
|
|
||||||
|
assert live["peak"] <= espn_dates.ESPN_CHUNK_WORKERS
|
||||||
|
assert live["peak"] > 1, "chunks should actually overlap"
|
||||||
|
|||||||
Reference in New Issue
Block a user