Skip to content

Commit 545d3ab

Browse files
jirhikerclaude
andcommitted
Dedup connectors and cut pipeline compute time
Deduplication: - Extract shared USGS retry/backoff and truncation checks into a composed USGSRequester (dependency-injected http client + warn), removing two near-identical inline retry blocks in the USGS connector. - Add a default __repr__ to BaseSource and remove 21 trivial per-connector overrides that only returned the class name. Performance: - Fetch a source's site chunks concurrently in _site_wrapper via a ThreadPoolExecutor, processing results in chunk order so output, site_limit, and rollback stay deterministic. Gated by the new config.fetch_workers knob (default 4; set 1 for legacy serial). - Tune httpx.Client connection pooling (keepalive reuse). - Stop sleeping after the final retry attempt in the base and USGS retry loops; switch USGS from linear to capped exponential backoff. - Fail fast on non-retryable client errors (4xx) instead of burning the full backoff schedule. All 360 tests pass. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 0d7cbed commit 545d3ab

11 files changed

Lines changed: 195 additions & 220 deletions

File tree

backend/config.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,12 @@ class Config:
152152
site_limit: int = 0
153153
dry: bool = False
154154

155+
# Number of chunks fetched concurrently per source (network-bound I/O).
156+
# 1 = serial (legacy behavior). Higher values speed up multi-chunk sources
157+
# but issue more simultaneous requests, which can trip per-source API rate
158+
# limits (e.g. USGS) — tune down if you see 429s.
159+
fetch_workers: int = 4
160+
155161
# date
156162
start_date: str = ""
157163
end_date: str = ""

backend/connectors/bor/source.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -38,9 +38,6 @@ class BORSiteSource(BaseSiteSource):
3838
def __init__(self):
3939
super().__init__(transformer=BORSiteTransformer())
4040

41-
def __repr__(self):
42-
return "BORSiteSource"
43-
4441
def health(self):
4542
try:
4643
resp = self.get_records()
@@ -66,9 +63,6 @@ def __init__(self):
6663
super().__init__(transformer=BORAnalyteTransformer())
6764
_source_parameter_name = None
6865

69-
def __repr__(self):
70-
return "BORAnalyteSource"
71-
7266
def _extract_parameter_record(self, record):
7367
record[PARAMETER_NAME] = self.config.parameter
7468
record[PARAMETER_VALUE] = record["attributes"]["result"]

backend/connectors/ckan/source.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -98,9 +98,6 @@ def __init__(self, resource_id, **kw):
9898
elif resource_id == ROSWELL_RESOURCE_ID:
9999
self.bounding_polygon = OSE_ROSWELL_ROSWELL_BOUNDING_POLYGON
100100

101-
def __repr__(self):
102-
return "NMOSERoswellSiteSource"
103-
104101
def health(self):
105102
params = self._get_params()
106103
params["limit"] = 1
@@ -123,9 +120,6 @@ def __init__(self, resource_id=None, **kw):
123120
kw.setdefault("transformer", OSERoswellWaterLevelTransformer())
124121
super().__init__(resource_id, **kw)
125122

126-
def __repr__(self):
127-
return "NMOSERoswellWaterLevelSource"
128-
129123
def get_records(self, site_record):
130124
return self._parse_response(site_record, self.get_response())
131125

backend/connectors/isc_seven_rivers/source.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -69,9 +69,6 @@ class ISCSevenRiversSiteSource(BaseSiteSource):
6969
def __init__(self):
7070
super().__init__(transformer=ISCSevenRiversSiteTransformer())
7171

72-
def __repr__(self):
73-
return "ISCSevenRiversSiteSource"
74-
7572
def health(self):
7673
try:
7774
resp = self.get_records()
@@ -94,9 +91,6 @@ class ISCSevenRiversAnalyteSource(BaseAnalyteSource):
9491
def __init__(self):
9592
super().__init__(transformer=ISCSevenRiversAnalyteTransformer())
9693

97-
def __repr__(self):
98-
return "ISCSevenRiversAnalyteSource"
99-
10094
def _get_analyte_id_and_name(self, analyte):
10195
""" """
10296
if self._analyte_ids is None:

backend/connectors/nmbgmr/source.py

Lines changed: 0 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,6 @@ class NMBGMRSiteSource(BaseSiteSource):
6464
def __init__(self):
6565
super().__init__(transformer=NMBGMRSiteTransformer())
6666

67-
def __repr__(self):
68-
return "NMBGMRSiteSource"
69-
7067
def health(self):
7168
try:
7269
resp = self._execute_json_request(
@@ -122,9 +119,6 @@ class NMBGMRAnalyteSource(BaseAnalyteSource):
122119
def __init__(self):
123120
super().__init__(transformer=NMBGMRAnalyteTransformer())
124121

125-
def __repr__(self):
126-
return "NMBGMRAnalyteSource"
127-
128122
def get_records(self, site_record):
129123
analyte = get_analyte_search_param(
130124
self.config.parameter, NMBGMR_ANALYTE_MAPPING
@@ -185,9 +179,6 @@ class NMBGMRWaterLevelSource(BaseWaterLevelSource):
185179
def __init__(self):
186180
super().__init__(transformer=NMBGMRWaterLevelTransformer())
187181

188-
def __repr__(self):
189-
return "NMBGMRWaterLevelSource"
190-
191182
def _clean_records(self, records):
192183
# remove records with no depth to water value
193184
return [

backend/connectors/nmenv/source.py

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -34,17 +34,13 @@
3434
URL = "https://nmenv.newmexicowaterdata.org/FROST-Server/v1.1/"
3535

3636

37-
3837
class DWBSiteSource(STSiteSource):
3938
url = URL
4039
bounding_polygon = NM_STATE_BOUNDING_POLYGON
4140

4241
def __init__(self):
4342
super().__init__(transformer=DWBSiteTransformer())
4443

45-
def __repr__(self):
46-
return "DWBSiteSource"
47-
4844
def health(self):
4945
try:
5046
resp = self.get_records(top=10, analyte=TDS)
@@ -115,9 +111,6 @@ class DWBAnalyteSource(STAnalyteSource):
115111
def __init__(self):
116112
super().__init__(transformer=DWBAnalyteTransformer())
117113

118-
def __repr__(self):
119-
return "DWBAnalyteSource"
120-
121114
def _parse_result(
122115
self, result, result_dt=None, result_id=None, result_location=None
123116
):

backend/connectors/st2/source.py

Lines changed: 0 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -184,48 +184,33 @@ class NMOSERoswellWaterLevelSource(ST2WaterLevelSource):
184184
def __init__(self):
185185
super().__init__(transformer=NMOSERoswellWaterLevelTransformer())
186186

187-
def __repr__(self):
188-
return "NMOSERoswellWaterLevelSource"
189-
190187

191188
class PVACDWaterLevelSource(ST2WaterLevelSource):
192189
agency = "PVACD"
193190

194191
def __init__(self):
195192
super().__init__(transformer=PVACDWaterLevelTransformer())
196193

197-
def __repr__(self):
198-
return "PVACDWaterLevelSource"
199-
200194

201195
class EBIDWaterLevelSource(ST2WaterLevelSource):
202196
agency = "EBID"
203197

204198
def __init__(self):
205199
super().__init__(transformer=EBIDWaterLevelTransformer())
206200

207-
def __repr__(self):
208-
return "EBIDWaterLevelSource"
209-
210201

211202
class BernCoWaterLevelSource(ST2WaterLevelSource):
212203
agency = "BernCo"
213204

214205
def __init__(self):
215206
super().__init__(transformer=BernCoWaterLevelTransformer())
216207

217-
def __repr__(self):
218-
return "BernCoWaterLevelSource"
219-
220208

221209
class CABQWaterLevelSource(ST2WaterLevelSource):
222210
agency = "CABQ"
223211

224212
def __init__(self):
225213
super().__init__(transformer=CABQWaterLevelTransformer())
226214

227-
def __repr__(self):
228-
return "CABQWaterLevelSource"
229-
230215

231216
# ============= EOF =============================================

0 commit comments

Comments
 (0)