Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 47 additions & 11 deletions statvar_imports/statistics_poland/download_input_data.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
from datetime import datetime
from google.cloud import storage
import io
from urllib3.util import Retry
from requests.adapters import HTTPAdapter

# Configure logging
logging.basicConfig(level=logging.INFO, format='%(levelname)s: %(message)s')
Expand Down Expand Up @@ -70,28 +72,58 @@ def get_template_map(template_df):
name_to_code[clean_name] = str(code).strip()
return name_to_code

def fetch_variables():
"""Fetches all variables for Subject P3447."""
def get_http_session(retries=10, backoff_factor=1.5):
"""Creates a requests session with automatic HTTP retries and backoff."""
session = requests.Session()
retry_strategy = Retry(
total=retries,
connect=retries,
read=retries,
status=retries,
backoff_factor=backoff_factor,
status_forcelist=[429, 500, 502, 503, 504],
allowed_methods=["GET"],
raise_on_status=False
)
adapter = HTTPAdapter(max_retries=retry_strategy)
session.mount("https://", adapter)
session.mount("http://", adapter)
return session

def make_request(session, url, headers=None, params=None, timeout=60):
"""Makes an HTTP GET request using the session's configured retry strategy and timeout."""
try:
return session.get(url, headers=headers, params=params, timeout=timeout)
except Exception as e:
logging.error(f"HTTP GET failed for {url}: {e}")
return None

def fetch_variables(session):
"""Fetches all variables for Subject P3447 with retries."""
logging.info(f"Downloading variable list for Subject {SUBJECT_ID}...")
v_map = {}

for page in range(10):
url = f"{API_BASE_URL}/variables?subject-id={SUBJECT_ID}&page-size=100&lang=pl&page={page}"
try:
resp = requests.get(url, headers=HEADERS, timeout=20)
if resp.status_code != 200: break
resp = make_request(session, url, headers=HEADERS, timeout=60)
if resp is None or resp.status_code != 200:
logging.warning(f"Metadata page {page} returned status {resp.status_code if resp else 'None'}")
break
data = resp.json()
results = data.get('results', [])
if not results: break
if not results:
break

for item in results:
full_name_parts = [str(v) for k, v in item.items() if k.startswith('n') and v]
full_name = " ".join(full_name_parts).lower()
v_map[str(item['id'])] = full_name

if len(results) < 100: break
if len(results) < 100:
break
except Exception as e:
logging.error(f"Metadata error page {page}: {e}")
logging.error(f"Metadata error page {page} after retries: {e}")
break

logging.info(f"Indexed {len(v_map)} variables.")
Expand All @@ -110,8 +142,9 @@ def download_and_process():
template_df.index.levels[1].astype(str)
])

session = get_http_session()
region_map = get_template_map(template_df)
v_metadata = fetch_variables()
v_metadata = fetch_variables(session)
if not v_metadata:
raise ValueError("Variable metadata failed to download.")

Expand Down Expand Up @@ -160,8 +193,10 @@ def download_and_process():
params.append(('year', str(y)))

try:
resp = requests.get(api_url, headers=HEADERS, params=params, timeout=20)
if resp.status_code != 200: continue
resp = make_request(session, api_url, headers=HEADERS, params=params, timeout=60)
if resp is None or resp.status_code != 200:
logging.error(f"Download returned status {resp.status_code if resp else 'None'} for var {var_id} level {lv}")
continue
results = resp.json().get('results', [])
if not results: continue

Expand Down Expand Up @@ -193,7 +228,8 @@ def download_and_process():
})
except Exception as e:
logging.error(f"Download Error on {var_id}: {e}")
time.sleep(0.05)
time.sleep(0.1)
time.sleep(0.1)

if not master_data:
raise ValueError("No data collected during the download loop.")
Expand Down
Loading