Files
Climate-Mood-Analysis/scripts/fetch_nsrdb_representative_point_ghi.py
2026-06-11 14:50:33 -04:00

544 lines
20 KiB
Python

#!/usr/bin/env python3
"""
Fetch NSRDB GOES TMY GHI data for county representative points.
This is designed for a small first test batch. It prompts for email and API key
when those values are not passed as command-line arguments.
"""
from __future__ import annotations
import argparse
import concurrent.futures
import csv
import getpass
import re
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from pathlib import Path
DEFAULT_POINTS_CSV = Path("data/nrel/county_representative_points.csv")
DEFAULT_OUTPUT_CSV = Path("data/nrel/county_representative_point_ghi_summary.csv")
DEFAULT_ERROR_CSV = Path("data/nrel/county_representative_point_ghi_error_log.csv")
DEFAULT_CACHE_DIR = Path("data/nrel/representative_point_csv")
DEFAULT_ENDPOINT = "https://developer.nlr.gov/api/nsrdb/v2/solar/nsrdb-GOES-tmy-v4-0-0-download.csv"
DEFAULT_POLAR_ENDPOINT = "https://developer.nlr.gov/api/nsrdb/v2/solar/nsrdb-polar-tmy-v4-0-0-download.csv"
DEFAULT_TMY_NAME = "tmy-2024"
DEFAULT_POLAR_TMY_NAME = "tmy"
DEFAULT_POLAR_MIN_LATITUDE = 60.0
EMAIL_PATTERN = re.compile(r"[\w.!#$%&'*+/=?^`{|}~-]+@[\w-]+(?:\.[\w-]+)+")
SENSITIVE_QUERY_PATTERN = re.compile(r"((?:api_key|email)=)([^&\s,\"']+)", re.IGNORECASE)
def prompt_for_secret(prompt: str, current_value: str | None) -> str:
"""Prompt for a sensitive value only when it was not supplied."""
if current_value:
return current_value
return getpass.getpass(prompt).strip()
def prompt_for_text(prompt: str, current_value: str | None) -> str:
"""Prompt for a normal text value only when it was not supplied."""
if current_value:
return current_value
return input(prompt).strip()
def read_county_points(points_csv: Path, start: int, limit: int | None) -> list[dict[str, str]]:
"""Read the county point CSV and return the requested batch."""
with points_csv.open(newline="", encoding="utf-8") as handle:
rows = list(csv.DictReader(handle))
if start < 0:
raise ValueError("--start must be 0 or greater.")
if limit is None:
return rows[start:]
if limit < 1:
raise ValueError("--limit must be 1 or greater.")
return rows[start : start + limit]
def read_previous_error_fips(error_csv: Path) -> set[str]:
"""Read county FIPS values from a previous error log."""
if not error_csv.exists():
return set()
with error_csv.open(newline="", encoding="utf-8") as handle:
return {
row["county_fips"]
for row in csv.DictReader(handle)
if row.get("county_fips")
}
def is_polar_candidate(point: dict[str, str], min_latitude: float) -> bool:
"""Return whether a point is in the latitude range for the Polar NSRDB endpoint."""
try:
return float(point["lat"]) >= min_latitude
except (KeyError, ValueError):
return False
def build_nsrdb_url(endpoint: str, api_key: str, email: str, point: dict[str, str], name: str) -> str:
"""Build a single-point NSRDB CSV request URL."""
wkt = f"POINT({point['lon']} {point['lat']})"
query = {
"api_key": api_key,
"wkt": wkt,
"attributes": "ghi",
"names": name,
"utc": "false",
"leap_day": "false",
"interval": "60",
"email": email,
}
return f"{endpoint}?{urllib.parse.urlencode(query)}"
def redact_url(url: str) -> str:
"""Return a request URL with sensitive query values removed."""
parsed_url = urllib.parse.urlsplit(url)
query = urllib.parse.parse_qsl(parsed_url.query, keep_blank_values=True)
redacted_query = [
(key, "<redacted>" if key in {"api_key", "email"} else value)
for key, value in query
]
return urllib.parse.urlunsplit(
parsed_url._replace(query=urllib.parse.urlencode(redacted_query))
)
def redact_sensitive_text(text: str) -> str:
"""Remove likely credentials and email addresses from log text."""
without_query_values = SENSITIVE_QUERY_PATTERN.sub(r"\1<redacted>", text)
return EMAIL_PATTERN.sub("<redacted>", without_query_values)
class NsrdDataRequestError(RuntimeError):
"""Store a failed NSRDB request with sanitized context for logging."""
def __init__(self, message: str, status_code: int | None = None, url: str | None = None, response: str = "") -> None:
super().__init__(message)
self.status_code = status_code
self.url = url
self.response = response
class RequestRateLimiter:
"""Coordinate network request starts across worker threads."""
def __init__(self, delay: float) -> None:
self.delay = max(0.0, delay)
self._lock = threading.Lock()
self._next_request_at = 0.0
def wait(self) -> None:
"""Wait until the next request can start."""
if self.delay <= 0:
return
with self._lock:
now = time.monotonic()
wait_seconds = max(0.0, self._next_request_at - now)
self._next_request_at = max(now, self._next_request_at) + self.delay
if wait_seconds:
time.sleep(wait_seconds)
def build_http_request_error(error: urllib.error.HTTPError, url: str) -> NsrdDataRequestError:
"""Create a structured NSRDB request error from an HTTPError."""
body = error.read().decode("utf-8-sig", errors="replace").strip()
body_excerpt = redact_sensitive_text(body[:1000]) if body else "No response body returned."
redacted_url = redact_url(url)
retry_after = error.headers.get("Retry-After") if error.headers else None
retry_message = f"\n Retry-After: {retry_after}" if retry_after else ""
message = f"HTTP {error.code} {error.reason}\n URL: {redacted_url}{retry_message}\n Response: {body_excerpt}"
return NsrdDataRequestError(
message=message,
status_code=error.code,
url=redacted_url,
response=body_excerpt,
)
def download_text(url: str, timeout: int, rate_limiter: RequestRateLimiter | None = None) -> str:
"""Download text from a URL and return its decoded body."""
request = urllib.request.Request(url, headers={"User-Agent": "county-climate-explorer/0.1"})
try:
if rate_limiter is not None:
rate_limiter.wait()
with urllib.request.urlopen(request, timeout=timeout) as response:
return response.read().decode("utf-8-sig")
except urllib.error.HTTPError as error:
raise build_http_request_error(error, url) from error
def read_or_download_county_csv(
point: dict[str, str],
endpoint: str,
api_key: str,
email: str,
name: str,
source_key: str,
cache_dir: Path,
timeout: int,
overwrite: bool,
rate_limiter: RequestRateLimiter | None = None,
) -> tuple[str, Path]:
"""Return cached NSRDB CSV text, downloading it first when needed."""
cache_dir.mkdir(parents=True, exist_ok=True)
cache_path = cache_dir / f"{point['county_fips']}_{source_key}_{name}_ghi.csv"
if cache_path.exists() and not overwrite:
return cache_path.read_text(encoding="utf-8-sig"), cache_path
url = build_nsrdb_url(endpoint, api_key, email, point, name)
csv_text = download_text(url, timeout, rate_limiter)
cache_path.write_text(csv_text, encoding="utf-8")
return csv_text, cache_path
def extract_ghi_values(csv_text: str) -> list[float]:
"""Extract numeric GHI values from an NSRDB CSV response."""
rows = list(csv.reader(csv_text.splitlines()))
header_index = next(
(
index
for index, row in enumerate(rows)
if {"Year", "Month", "Day", "Hour", "Minute", "GHI"}.issubset(set(row))
),
None,
)
if header_index is None:
raise ValueError("Could not find the NSRDB data header row with a GHI column.")
ghi_index = rows[header_index].index("GHI")
ghi_values: list[float] = []
for row in rows[header_index + 1 :]:
if len(row) <= ghi_index or not row[ghi_index].strip():
continue
ghi_values.append(float(row[ghi_index]))
if not ghi_values:
raise ValueError("No GHI values were found in the NSRDB response.")
return ghi_values
def summarize_ghi(point: dict[str, str], csv_text: str, source_file: Path, source_label: str) -> dict[str, str]:
"""Convert hourly GHI values into average daily kWh/m2/day."""
ghi_values = extract_ghi_values(csv_text)
ghi_sum = sum(ghi_values)
avg_daily_ghi = ghi_sum / 1000 / 365
return {
"county_fips": point["county_fips"],
"county_name": point["county_name"],
"state_fips": point["state_fips"],
"state_abbr": point["state_abbr"],
"lat": point["lat"],
"lon": point["lon"],
"avgSolarGhiKwhM2Day": f"{avg_daily_ghi:.3f}",
"ghi_rows": str(len(ghi_values)),
"ghi_min": f"{min(ghi_values):.1f}",
"ghi_max": f"{max(ghi_values):.1f}",
"source": f"{source_label} representative point",
"raw_csv": str(source_file),
}
def write_summary_csv(rows: list[dict[str, str]], output_csv: Path) -> None:
"""Write fetched county solar summaries to CSV."""
output_csv.parent.mkdir(parents=True, exist_ok=True)
with output_csv.open("w", newline="", encoding="utf-8") as handle:
writer = csv.DictWriter(
handle,
fieldnames=[
"county_fips",
"county_name",
"state_fips",
"state_abbr",
"lat",
"lon",
"avgSolarGhiKwhM2Day",
"ghi_rows",
"ghi_min",
"ghi_max",
"source",
"raw_csv",
],
)
writer.writeheader()
writer.writerows(rows)
def write_error_csv(rows: list[dict[str, str]], error_csv: Path) -> None:
"""Write failed county solar fetches to CSV."""
error_csv.parent.mkdir(parents=True, exist_ok=True)
with error_csv.open("w", newline="", encoding="utf-8") as handle:
writer = csv.DictWriter(
handle,
fieldnames=[
"county_fips",
"county_name",
"state_fips",
"state_abbr",
"lat",
"lon",
"error_type",
"status_code",
"request_url",
"response",
"message",
],
)
writer.writeheader()
writer.writerows(rows)
def build_error_row(point: dict[str, str], error: Exception) -> dict[str, str]:
"""Convert a failed county fetch into a structured CSV row."""
status_code = ""
request_url = ""
response = ""
if isinstance(error, NsrdDataRequestError):
status_code = str(error.status_code or "")
request_url = error.url or ""
response = error.response
return {
"county_fips": point["county_fips"],
"county_name": point["county_name"],
"state_fips": point["state_fips"],
"state_abbr": point["state_abbr"],
"lat": point["lat"],
"lon": point["lon"],
"error_type": type(error).__name__,
"status_code": status_code,
"request_url": request_url,
"response": redact_sensitive_text(response),
"message": redact_sensitive_text(str(error)),
}
def build_request_profiles(
args: argparse.Namespace,
point: dict[str, str],
previous_error_fips: set[str],
) -> list[dict[str, str]]:
"""Build the ordered list of NSRDB endpoints to try for a county point."""
goes_profile = {
"source_key": "goes-tmy",
"source_label": "NSRDB GOES TMY PSM v4",
"endpoint": args.endpoint,
"name": args.name,
}
polar_profile = {
"source_key": "polar-tmy",
"source_label": "NSRDB Polar TMY PSM v4",
"endpoint": args.polar_endpoint,
"name": args.polar_name,
}
if not args.polar_fallback or not is_polar_candidate(point, args.polar_min_latitude):
return [goes_profile]
if point["county_fips"] in previous_error_fips:
return [polar_profile, goes_profile]
return [goes_profile, polar_profile]
def should_try_next_profile(error: Exception, profile_index: int, profiles: list[dict[str, str]]) -> bool:
"""Return whether another configured endpoint should be tried after this error."""
if profile_index >= len(profiles) - 1:
return False
if not isinstance(error, NsrdDataRequestError):
return False
return error.status_code == 400 and "No data available at the provided location" in error.response
def fetch_county_summary(
args: argparse.Namespace,
point: dict[str, str],
previous_error_fips: set[str],
api_key: str,
email: str,
rate_limiter: RequestRateLimiter,
) -> tuple[dict[str, str] | None, dict[str, str] | None, list[str]]:
"""Fetch and summarize one county point."""
profiles = build_request_profiles(args, point, previous_error_fips)
final_error: Exception | None = None
messages: list[str] = []
for profile_index, profile in enumerate(profiles):
try:
messages.append(f"Trying {profile['source_label']} ({profile['name']})")
csv_text, source_file = read_or_download_county_csv(
point=point,
endpoint=profile["endpoint"],
api_key=api_key,
email=email,
name=profile["name"],
source_key=profile["source_key"],
cache_dir=args.cache_dir,
timeout=args.timeout,
overwrite=args.overwrite,
rate_limiter=rate_limiter,
)
summary = summarize_ghi(point, csv_text, source_file, profile["source_label"])
return summary, None, messages
except (OSError, urllib.error.URLError, RuntimeError, ValueError) as error:
final_error = error
if should_try_next_profile(error, profile_index, profiles):
messages.append(f"{profile['source_label']} had no data; trying fallback.")
continue
break
if final_error is None:
final_error = RuntimeError("No NSRDB profiles were available for this county.")
return None, build_error_row(point, final_error), messages
def log_county_result(
index: int,
total: int,
point: dict[str, str],
summary: dict[str, str] | None,
error: dict[str, str] | None,
messages: list[str],
) -> None:
"""Print progress for one completed county."""
label = f"{point['county_fips']} {point['county_name']}, {point['state_abbr']}"
status = "Fetched" if summary is not None else "Skipped"
print(f"[{index}/{total}] {status} {label}")
for message in messages:
print(f" {message}")
if error is not None:
print(f" Error: {error['message']}")
def fetch_batch(args: argparse.Namespace) -> list[dict[str, str]]:
"""Fetch and summarize the requested county point batch."""
email = prompt_for_text("NLR/NREL API email: ", args.email)
api_key = prompt_for_secret("NLR/NREL API key: ", args.api_key)
points = read_county_points(args.points_csv, args.start, args.limit)
previous_error_fips = read_previous_error_fips(args.error_csv)
rate_limiter = RequestRateLimiter(args.delay)
summaries_by_index: dict[int, dict[str, str]] = {}
errors_by_index: dict[int, dict[str, str]] = {}
total = len(points)
if args.workers == 1:
for index, point in enumerate(points, start=1):
summary, error, messages = fetch_county_summary(
args=args,
point=point,
previous_error_fips=previous_error_fips,
api_key=api_key,
email=email,
rate_limiter=rate_limiter,
)
if summary is not None:
summaries_by_index[index] = summary
if error is not None:
errors_by_index[index] = error
log_county_result(index, total, point, summary, error, messages)
else:
with concurrent.futures.ThreadPoolExecutor(max_workers=args.workers) as executor:
future_to_context = {
executor.submit(
fetch_county_summary,
args,
point,
previous_error_fips,
api_key,
email,
rate_limiter,
): (index, point)
for index, point in enumerate(points, start=1)
}
for future in concurrent.futures.as_completed(future_to_context):
index, point = future_to_context[future]
summary, error, messages = future.result()
if summary is not None:
summaries_by_index[index] = summary
if error is not None:
errors_by_index[index] = error
log_county_result(index, total, point, summary, error, messages)
errors = [errors_by_index[index] for index in sorted(errors_by_index)]
write_error_csv(errors, args.error_csv)
return [summaries_by_index[index] for index in sorted(summaries_by_index)]
def parse_args() -> argparse.Namespace:
"""Parse command-line arguments."""
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--points-csv", type=Path, default=DEFAULT_POINTS_CSV)
parser.add_argument("--output-csv", type=Path, default=DEFAULT_OUTPUT_CSV)
parser.add_argument("--error-csv", type=Path, default=DEFAULT_ERROR_CSV)
parser.add_argument("--representative-point-csv-dir", "--cache-dir", dest="cache_dir", type=Path, default=DEFAULT_CACHE_DIR)
parser.add_argument("--endpoint", default=DEFAULT_ENDPOINT)
parser.add_argument("--polar-endpoint", default=DEFAULT_POLAR_ENDPOINT)
parser.add_argument("--name", default=DEFAULT_TMY_NAME, help="NSRDB TMY name, such as tmy or tmy-2024.")
parser.add_argument("--polar-name", default=DEFAULT_POLAR_TMY_NAME, help="NSRDB Polar TMY name, usually tmy.")
parser.add_argument("--polar-min-latitude", type=float, default=DEFAULT_POLAR_MIN_LATITUDE)
parser.add_argument(
"--no-polar-fallback",
action="store_false",
dest="polar_fallback",
help="Disable retrying failed high-latitude GOES requests against the NSRDB Polar endpoint.",
)
parser.add_argument("--email", help="Email address registered with the NLR/NREL API.")
parser.add_argument("--api-key", help="API key. If omitted, the script prompts securely.")
parser.add_argument("--start", type=int, default=0, help="Zero-based row offset in the points CSV.")
parser.add_argument("--limit", type=int, default=10, help="Number of counties to fetch for the test batch.")
parser.add_argument(
"--delay",
type=float,
default=1.0,
help="Minimum seconds between new API requests. Cached files do not wait.",
)
parser.add_argument(
"--workers",
type=int,
default=4,
help="Number of county jobs to run at once. Set to 1 for the old sequential behavior.",
)
parser.add_argument("--timeout", type=int, default=120, help="Request timeout in seconds.")
parser.add_argument("--overwrite", action="store_true", help="Re-download cached county CSV files.")
args = parser.parse_args()
if args.delay < 0:
parser.error("--delay must be 0 or greater.")
if args.workers < 1:
parser.error("--workers must be 1 or greater.")
return args
def main() -> None:
"""Run the NSRDB point-fetch test batch."""
args = parse_args()
summaries = fetch_batch(args)
write_summary_csv(summaries, args.output_csv)
print(f"Wrote {len(summaries)} solar summaries to {args.output_csv}")
if __name__ == "__main__":
main()