diff --git a/cwms/timeseries/timeseries.py b/cwms/timeseries/timeseries.py index 894989f..33d5fb1 100644 --- a/cwms/timeseries/timeseries.py +++ b/cwms/timeseries/timeseries.py @@ -1,5 +1,6 @@ import concurrent.futures import logging +import re from datetime import datetime, timedelta, timezone from typing import Any, Dict, List, Optional, Tuple @@ -10,6 +11,74 @@ from cwms.catalog.catalog import get_ts_extents from cwms.cwms_types import JSON, Data +_DEFAULT_CHUNK_DAYS = 365 +_MIN_INTERVAL_MINUTES = 2 +_FINE_INTERVAL_MINUTES = 15 +_FINE_INTERVAL_CHUNK_DAYS = 365 +_HOURLY_CHUNK_DAYS = 365 +_SIX_HOURLY_CHUNK_DAYS = 1460 +_COARSE_INTERVAL_CHUNK_DAYS = 2920 +_INTERVAL_PATTERN = re.compile( + r"^(?P\d+)(?P" + r"Minute|Minutes|Hour|Hours|Day|Days|" + r"Week|Weeks|Month|Months|Year|Years)$" +) +_INTERVAL_MINUTES = { + "Minute": 1, + "Minutes": 1, + "Hour": 60, + "Hours": 60, + "Day": 24 * 60, + "Days": 24 * 60, + "Week": 7 * 24 * 60, + "Weeks": 7 * 24 * 60, + "Month": 30 * 24 * 60, + "Months": 30 * 24 * 60, + "Year": 365 * 24 * 60, + "Years": 365 * 24 * 60, +} + + +def get_timeseries_chunk_size(ts_id: str) -> timedelta: + """Return the default request chunk size for a time series interval. + + Local regular time series intervals, such as ``~15Minutes``, use the same + chunk size as their regular interval. Unrecognized intervals retain the + conservative default used for 15-minute through hourly data. + """ + ts_id_parts = ts_id.split(".") + if len(ts_id_parts) < 4: + return timedelta(days=_DEFAULT_CHUNK_DAYS) + + interval = ts_id_parts[3].removeprefix("~") + match = _INTERVAL_PATTERN.fullmatch(interval) + if match is None: + return timedelta(days=_DEFAULT_CHUNK_DAYS) + + interval_minutes = ( + int(match.group("count")) * _INTERVAL_MINUTES[match.group("unit")] + ) + if interval_minutes < _MIN_INTERVAL_MINUTES: + return timedelta(days=_DEFAULT_CHUNK_DAYS) + + # These bands balance request overhead and response size based on production + # CDA timings. Fine intervals scale toward about 35,000 expected values. + if interval_minutes < _FINE_INTERVAL_MINUTES: + chunk_days = max( + 1, + round( + _FINE_INTERVAL_CHUNK_DAYS * interval_minutes / _FINE_INTERVAL_MINUTES + ), + ) + elif interval_minutes <= 60: + chunk_days = _HOURLY_CHUNK_DAYS + elif interval_minutes <= 6 * 60: + chunk_days = _SIX_HOURLY_CHUNK_DAYS + else: + chunk_days = _COARSE_INTERVAL_CHUNK_DAYS + + return timedelta(days=chunk_days) + def get_multi_timeseries_df( ts_ids: list[str], @@ -136,6 +205,9 @@ def chunk_timeseries_time_range( List[Tuple[datetime, datetime]] A list of tuples, where each tuple represents the start and end of a chunk. """ + if chunk_size <= timedelta(0): + raise ValueError("chunk_size must be greater than zero") + chunks = [] current = begin while current < end: @@ -298,7 +370,7 @@ def get_timeseries( trim: Optional[bool] = True, multithread: Optional[bool] = True, max_workers: int = 20, - max_days_per_chunk: int = 14, + max_days_per_chunk: Optional[int] = None, ) -> Data: """Retrieves time series values from a specified time series and time window. Value date-times obtained are always in UTC. @@ -338,15 +410,18 @@ def get_timeseries( trim: boolean, optional, default is True Specifies whether to trim missing values from the beginning and end of the retrieved values. multithread: boolean, optional, default is True - Specifies whether to trim missing values from the beginning and end of the retrieved values. + Specifies whether to retrieve time series chunks concurrently. max_workers: integer, default is 20 - The maximum number of worker threads that will be spawned for multithreading, If calling more than 3 years of 15 minute data, consider using 30 max_workers - max_days_per_chunk: integer, default is 14 - The maximum number of days that would be included in a thread. If calling more than 1 year of 15 minute data, consider using 30 days + The maximum number of worker threads used for concurrent requests. + max_days_per_chunk: integer, optional, default is None + The maximum number of days included in each request. By default, + the chunk size is selected from the time series interval. Returns ------- cwms data type. data.json will return the JSON output and data.df will return a dataframe. dates are all in UTC """ + if max_days_per_chunk is not None and max_days_per_chunk <= 0: + raise ValueError("max_days_per_chunk must be greater than zero") selector = "values" endpoint = "timeseries" @@ -386,8 +461,12 @@ def get_timeseries( ) return Data(response, selector=selector) - # divide the time range into chunks - chunks = chunk_timeseries_time_range(begin, end, timedelta(days=max_days_per_chunk)) + chunk_size = ( + timedelta(days=max_days_per_chunk) + if max_days_per_chunk is not None + else get_timeseries_chunk_size(ts_id) + ) + chunks = chunk_timeseries_time_range(begin, end, chunk_size) # find max worker thread max_workers = max(min(len(chunks), max_workers), 1) diff --git a/tests/mock/timeseries/timeseries_test.py b/tests/mock/timeseries/timeseries_test.py index e5c906b..911e77c 100644 --- a/tests/mock/timeseries/timeseries_test.py +++ b/tests/mock/timeseries/timeseries_test.py @@ -3,7 +3,7 @@ # All Rights Reserved. USACE PROPRIETARY/CONFIDENTIAL. # Source may not be released without written approval from HEC -from datetime import datetime +from datetime import datetime, timedelta import pandas as pd import pytest @@ -272,6 +272,128 @@ def failing_call(): assert call_count == 1 +@pytest.mark.parametrize( + ("ts_id", "expected_days"), + [ + ("Test.Stage.Inst.2Minutes.0.Test", 49), + ("Test.Stage.Inst.5Minutes.0.Test", 122), + ("Test.Stage.Inst.15Minutes.0.Test", 365), + ("Test.Stage.Inst.~15Minutes.0.Test", 365), + ("Test.Stage.Inst.1Hour.0.Test", 365), + ("Test.Stage.Inst.~1Hour.0.Test", 365), + ("Test.Stage.Inst.6Hours.0.Test", 1460), + ("Test.Stage.Inst.~6Hours.0.Test", 1460), + ("Test.Stage.Inst.1Day.0.Test", 2920), + ("Test.Stage.Inst.~1Day.0.Test", 2920), + ("Test.Stage.Inst.1Week.0.Test", 2920), + ("Test.Stage.Inst.1Month.0.Test", 2920), + ("Test.Stage.Inst.1Year.0.Test", 2920), + ], +) +def test_get_timeseries_chunk_size(ts_id, expected_days): + assert timeseries.get_timeseries_chunk_size(ts_id) == timedelta(days=expected_days) + + +@pytest.mark.parametrize( + "ts_id", + [ + "Invalid", + "Test.Stage.Inst.0.0.Test", + "Test.Stage.Inst.1Second.0.Test", + "Test.Stage.Inst.1Minute.0.Test", + "Test.Stage.Inst.Irregular.0.Test", + ], +) +def test_get_timeseries_chunk_size_falls_back_for_unknown_intervals(ts_id): + assert timeseries.get_timeseries_chunk_size(ts_id) == timedelta(days=365) + + +def test_chunk_timeseries_time_range_rejects_nonpositive_size(): + now = datetime.now(tz=pytz.UTC) + + with pytest.raises(ValueError, match="chunk_size must be greater than zero"): + timeseries.chunk_timeseries_time_range( + now, now + timedelta(days=1), timedelta() + ) + + +@pytest.mark.parametrize( + ("ts_id", "expected_days"), + [ + ("Test.Stage.Inst.~15Minutes.0.DefaultChunk", 365), + ("Test.Stage.Inst.1Day.0.DefaultChunk", 2920), + ], +) +def test_get_timeseries_uses_interval_chunk_size_by_default( + monkeypatch, ts_id, expected_days +): + begin = datetime(2025, 1, 1, tzinfo=pytz.UTC) + end = begin + timedelta(days=expected_days * 2 + 1) + expected = object() + captured = {} + + def fetch_chunks(chunks, params, selector, endpoint, max_workers): + captured["chunks"] = chunks + captured["max_workers"] = max_workers + return [object()] + + monkeypatch.setattr(timeseries, "fetch_timeseries_chunks", fetch_chunks) + monkeypatch.setattr( + timeseries, "combine_timeseries_results", lambda results: expected + ) + + result = timeseries.get_timeseries( + ts_id=ts_id, + office_id="MVP", + begin=begin, + end=end, + ) + + assert result is expected + assert captured["chunks"] == timeseries.chunk_timeseries_time_range( + begin, end, timedelta(days=expected_days) + ) + assert captured["max_workers"] == 3 + + +def test_get_timeseries_max_days_per_chunk_overrides_interval(monkeypatch): + begin = datetime(2025, 1, 1, tzinfo=pytz.UTC) + end = begin + timedelta(days=31) + captured = {} + + def fetch_chunks(chunks, params, selector, endpoint, max_workers): + captured["chunks"] = chunks + return [object()] + + monkeypatch.setattr(timeseries, "fetch_timeseries_chunks", fetch_chunks) + monkeypatch.setattr( + timeseries, "combine_timeseries_results", lambda results: object() + ) + + timeseries.get_timeseries( + ts_id="Test.Stage.Inst.1Day.0.OverrideChunk", + office_id="MVP", + begin=begin, + end=end, + max_days_per_chunk=10, + ) + + assert captured["chunks"] == timeseries.chunk_timeseries_time_range( + begin, end, timedelta(days=10) + ) + + +def test_get_timeseries_rejects_nonpositive_max_days_per_chunk(): + with pytest.raises( + ValueError, match="max_days_per_chunk must be greater than zero" + ): + timeseries.get_timeseries( + ts_id="Test.Stage.Inst.15Minutes.0.InvalidChunk", + office_id="MVP", + max_days_per_chunk=0, + ) + + def test_get_timeseries_group_default(requests_mock): group_id = "USGS TS Data Acquisition" category_id = "Data Acquisition"