Загрузка данных
import builtins
import json
import os
import sys
import logging
import traceback
import requests
from datetime import datetime, timedelta
from calendar import monthrange
from argparse import ArgumentParser
from io import StringIO
import csv
import time
from umdap_config import LOG_FOLDER, ENV_FOLDER, EROUTER_HOST, PRIMARY_MADAB, UMDAP_ETL_URL_LIST
APPLICATION_NAME = 'externalGTT'
LOG_FILE = os.path.join(LOG_FOLDER, "externalGTT.log")
CONFIG_FILE = os.path.join(ENV_FOLDER, "externalGtt.json")
CREDENTIALS_FILE = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.json")
BASE_URL = 'https://www.globaltradetracker.com/api/rest'
builtins.BASE_URL = BASE_URL
builtins.XML_DATE_FMT = '%d.%m.%Y %H:%M:%S'
builtins.DELIMITER = ','
results = list()
last_update_file = os.path.join(os.path.dirname(os.path.abspath(__file__)), "last_update.json")
def log(msg=''):
"""Логирование с временной меткой"""
str_with_time = f'{datetime.now().strftime("%Y/%m/%d %H:%M:%S")} {msg}'
print(str_with_time)
with open(LOG_FILE, "a", encoding="utf-8") as file:
file.write(f'{str_with_time}\n')
builtins.log = log
def get_token():
"""Получение токена для API"""
try:
with open(CREDENTIALS_FILE, 'r') as f:
credentials = json.load(f)
link = f"{BASE_URL}/gettoken?userid={credentials['user_id']}&password={credentials['password']}"
log(f"Getting token from GTT API")
response = requests.get(link, timeout=30)
if response.status_code == 200:
token = response.content.decode().strip()
if token:
log("Token received successfully")
return token
else:
log("Empty token received")
return None
else:
log(f"Failed to get token. Status code: {response.status_code}")
log(f"Response: {response.text}")
return None
except FileNotFoundError:
log(f"ERROR: Credentials file not found: {CREDENTIALS_FILE}")
return None
except Exception as e:
log(f"Error getting token: {str(e)}")
return None
def get_subscription_countries(token):
"""Получение списка стран из подписки (только default reporters)"""
try:
link = f"{BASE_URL}/countries?token={token}"
log(f"Getting available countries from subscription")
response = requests.get(link, timeout=30)
if response.status_code == 200:
data = response.json()
countries = []
if isinstance(data, list):
for item in data:
if isinstance(item, dict):
country_data = item.get('country', item)
if country_data.get('reportingcountry') == 'true' and country_data.get('defaultReporter') == True:
code = country_data.get('reportercode', '')
if code:
countries.append(code)
countries = list(dict.fromkeys(countries))
# ВРЕМЕННО: только 9 стран для теста
test_countries = [
'US_district',
'BR_national',
'IN_national',
'CN_national',
'DE_national',
'RU_national',
'JP_national',
'KR_national',
'TR_national'
]
countries = [c for c in test_countries if c in countries]
log(f"Test countries ({len(countries)}): {countries}")
return countries
else:
log(f"ERROR: Failed to get countries. Status: {response.status_code}")
return []
except Exception as e:
log(f"ERROR: Failed to get countries: {str(e)}")
return []
def get_subscription_hs_codes(token):
"""Получение списка HS кодов из подписки (динамически)"""
try:
link = f"{BASE_URL}/subscriptions?token={token}"
log(f"Getting available HS codes from subscription")
response = requests.get(link, timeout=30)
if response.status_code == 200:
data = response.json()
hs_codes = []
if isinstance(data, list):
for item in data:
if isinstance(item, dict):
codes = item.get('hsCodes', [])
if codes:
hs_codes.extend([str(code) for code in codes])
# ВРЕМЕННО: только 1 HS-код для теста
hs_codes = ['31']
log(f"Test HS codes ({len(hs_codes)}): {hs_codes}")
return hs_codes
else:
log(f"ERROR: Failed to get subscriptions. Status: {response.status_code}")
return []
except Exception as e:
log(f"ERROR: Failed to get subscriptions: {str(e)}")
return []
def get_last_update_time():
"""Чтение времени последнего обновления"""
try:
if os.path.exists(last_update_file):
with open(last_update_file, 'r') as f:
data = json.load(f)
return data.get('last_update', None)
return None
except:
return None
def save_last_update_time(update_time):
"""Сохранение времени последнего обновления"""
try:
with open(last_update_file, 'w') as f:
json.dump({'last_update': update_time}, f)
log(f"Last update time saved: {update_time}")
except Exception as e:
log(f"Error saving last update time: {str(e)}")
def check_data_updates(token):
"""Проверка наличия обновлений данных"""
try:
last_update = get_last_update_time()
if not last_update:
return True
link = f"{BASE_URL}/dataupdates?token={token}&updatedAfter={last_update}"
log(f"Checking for data updates since {last_update}")
response = requests.get(link, timeout=30)
if response.status_code == 200:
data = response.json()
if data and len(data) > 0:
log(f"Found {len(data)} updates available")
return True
else:
log("No updates found")
return False
else:
log(f"Error checking updates. Status code: {response.status_code}")
return False
except Exception as e:
log(f"Error checking updates: {str(e)}")
return False
def parse_gtt_data(data, reporter_code="", trade_flow=""):
"""Парсинг данных из JSON ответа GTT API"""
global results
if not data:
return
if isinstance(data, list):
for record in data:
parse_single_record(record, reporter_code, trade_flow)
elif isinstance(data, dict):
if 'data' in data:
parse_gtt_data(data['data'], reporter_code, trade_flow)
elif 'records' in data:
parse_gtt_data(data['records'], reporter_code, trade_flow)
elif 'result' in data:
parse_gtt_data(data['result'], reporter_code, trade_flow)
elif 'items' in data:
parse_gtt_data(data['items'], reporter_code, trade_flow)
else:
for key, value in data.items():
if isinstance(value, list):
parse_gtt_data(value, reporter_code, trade_flow)
return
parse_single_record(data, reporter_code, trade_flow)
def parse_single_record(record, reporter_code="", trade_flow=""):
"""Парсинг одной записи данных"""
try:
if not isinstance(record, dict):
return
def safe_get(obj, key, default=''):
if isinstance(obj, dict):
return obj.get(key, default)
return default
reporter = record.get('reporter', {})
partner = record.get('partner', {})
commodity = record.get('commodity', {})
period_data = record.get('period', None)
period_str = parse_period(period_data)
value = record.get('value', record.get('monetaryValue', 0))
if isinstance(value, dict):
value = value.get('number', value.get('value', 0))
quantity1 = record.get('quantity1', record.get('primaryQuantity', 0))
if isinstance(quantity1, dict):
quantity1 = quantity1.get('number', quantity1.get('value', 0))
quantity2 = record.get('quantity2', record.get('secondaryQuantity', 0))
if isinstance(quantity2, dict):
quantity2 = quantity2.get('number', quantity2.get('value', 0))
new_result = {
'Reporter_Code': safe_get(reporter, 'code', reporter_code) if isinstance(reporter, dict) else reporter_code,
'Trade_Flow': record.get('tradeFlow', record.get('flow', trade_flow)),
'Is_Mirror_Data': record.get('isMirrorData', ''),
'Period': period_str,
'Reporter_Name': safe_get(reporter, 'name', '') if isinstance(reporter, dict) else '',
'Reporter_Description': safe_get(reporter, 'description', '') if isinstance(reporter, dict) else '',
'Reporter_Source': safe_get(reporter, 'source', '') if isinstance(reporter, dict) else '',
'Incoterm': record.get('incoterm', ''),
'Partner_Code': safe_get(partner, 'code', '') if isinstance(partner, dict) else '',
'Partner_Name': safe_get(partner, 'name', '') if isinstance(partner, dict) else '',
'Commodity_HS_Code': safe_get(commodity, 'hsCode', record.get('hsCode', '')),
'Commodity_Description': safe_get(commodity, 'description', ''),
'HS6_Code': record.get('hs6Code', safe_get(commodity, 'hs6Code', '')),
'HS6_Code_Description': record.get('hs6Description', ''),
'Subdivision': record.get('subdivision', ''),
'Port': record.get('port', ''),
'Transport': record.get('transport', ''),
'Foreign_Port': record.get('foreignPort', ''),
'US_State': record.get('usState', ''),
'Customs_Regime': record.get('customsRegime', ''),
'Suppression': record.get('suppression', ''),
'Monetary_Value': value,
'Currency': record.get('currency', 'USD'),
'Secondary_Monetary_Value': record.get('secondaryValue', record.get('secondaryMonetaryValue', 0)),
'Secondary_Currency': record.get('secondaryCurrency', ''),
'Primary_Quantity': quantity1,
'Primary_Quantity_Unit': record.get('quantity1Unit', record.get('primaryQuantityUnit', '')),
'Secondary_Quantity': quantity2,
'Secondary_Quantity_Unit': record.get('quantity2Unit', record.get('secondaryQuantityUnit', '')),
'Primary_Quantity_Price': record.get('price1', record.get('primaryQuantityPrice', 0)),
'Primary_Quantity_Price_Unit': record.get('price1Unit', record.get('primaryQuantityPriceUnit', '')),
'Secondary_Quantity_Price': record.get('price2', record.get('secondaryQuantityPrice', 0)),
'Secondary_Quantity_Price_Unit': record.get('price2Unit', record.get('secondaryQuantityPriceUnit', '')),
}
results.append(new_result)
except Exception as e:
log(f"Error parsing record: {str(e)}")
log(f"Record data: {str(record)[:500]}")
def parse_period(period_data):
"""Парсинг периода в формат даты"""
try:
if isinstance(period_data, dict):
year = period_data.get('year', period_data.get(0, 2000))
month = period_data.get('month', period_data.get(1, 1))
elif isinstance(period_data, list) and len(period_data) >= 2:
year, month = period_data[0], period_data[1]
elif isinstance(period_data, str):
parts = period_data.split('-')
if len(parts) >= 2:
year, month = int(parts[0]), int(parts[1])
if len(parts) == 3:
day = int(parts[2])
return datetime(year, month, day).strftime('%Y-%m-%d')
else:
return datetime(2000, 1, 1).strftime('%Y-%m-%d')
else:
return datetime(2000, 1, 1).strftime('%Y-%m-%d')
last_day = monthrange(int(year), int(month))[1]
return datetime(int(year), int(month), last_day).strftime('%Y-%m-%d')
except Exception as e:
log(f"Error parsing period {period_data}: {str(e)}")
return datetime(2000, 1, 1).strftime('%Y-%m-%d')
def get_historical_data(token, countries, hs_codes):
"""БЛОК 1: Получение исторических данных только за 2026 год"""
log("=" * 60)
log("Starting TEST data collection for 2026")
log(f"Countries ({len(countries)}): {countries}")
log(f"HS Codes ({len(hs_codes)}): {hs_codes}")
log("=" * 60)
year = 2026
current_month = datetime.now().month
total_requests = 0
successful_requests = 0
empty_requests = 0
error_requests = 0
log(f"\n{'='*60}")
log(f"Processing year: {year}")
log(f"{'='*60}")
for country in countries:
for hscode in hs_codes:
for impexp in ['E', 'I']:
try:
from_date = f"{year}-01"
to_date = f"{year}-{current_month:02d}"
link = f"{BASE_URL}/getreport"
params = {
'token': token,
'hscode': hscode,
'reporter': country,
'impexp': impexp,
'from': from_date,
'to': to_date,
'currency': 'USD',
'format': 'json'
}
total_requests += 1
response = requests.get(link, params=params, timeout=300)
if response.status_code == 200:
successful_requests += 1
data = response.json()
if data and isinstance(data, list) and len(data) > 0:
parse_gtt_data(data, country, impexp)
log(f"✓ [{total_requests}] Country={country}, HS={hscode}, Flow={impexp}, Records={len(data)}, Total={len(results)}")
else:
empty_requests += 1
else:
error_requests += 1
log(f"✗ [{total_requests}] Country={country}, HS={hscode}, Flow={impexp}, Error={response.status_code}")
except Exception as e:
log(f"✗ Exception: {str(e)}")
continue
log(f"\n{'='*60}")
log(f"COMPLETED")
log(f"Total requests: {total_requests}")
log(f"Successful: {successful_requests}")
log(f"Empty: {empty_requests}")
log(f"Errors: {error_requests}")
log(f"Records collected: {len(results)}")
log(f"{'='*60}")
def get_updated_data(token, countries, hs_codes):
"""БЛОК 2: Получение обновленных данных"""
log("Checking for updated data")
last_update = get_last_update_time()
current_date = datetime.now().strftime('%Y-%m')
if not check_data_updates(token):
log("No updates available")
return
for country in countries:
for hscode in hs_codes:
for impexp in ['E', 'I']:
try:
link = f"{BASE_URL}/getreport"
params = {
'token': token,
'hscode': hscode,
'reporter': country,
'impexp': impexp,
'updatedAfter': last_update if last_update else '2000-01',
'latestavailablemonths': 3,
'currency': 'USD',
'format': 'json'
}
response = requests.get(link, params=params, timeout=300)
if response.status_code == 200:
data = response.json()
parse_gtt_data(data, country, impexp)
except Exception as e:
log(f"Error: {str(e)}")
continue
save_last_update_time(current_date)
log(f"Updated data collection completed. Total records: {len(results)}")
def gtt_request(url):
"""Формирование CSV для отправки в DWH"""
global results
log(f'{len(results)} import(s) generated')
if results:
with StringIO() as f:
writer = csv.DictWriter(f, results[0].keys())
writer.writeheader()
writer.writerows(results)
f.seek(0)
return f.read()
return ""
def main():
"""Основная функция"""
global results
log('--------------------- Script started ------------------------------')
try:
token = get_token()
if not token:
log("ERROR: Failed to get token")
sys.exit(1)
countries = get_subscription_countries(token)
hs_codes = get_subscription_hs_codes(token)
if not countries or not hs_codes:
log("ERROR: Failed to get subscription parameters")
sys.exit(1)
last_update = get_last_update_time()
if last_update is None:
log("First run detected - starting TEST data collection")
get_historical_data(token, countries, hs_codes)
save_last_update_time(datetime.now().strftime('%Y-%m'))
else:
log(f"Regular update - last update was {last_update}")
get_updated_data(token, countries, hs_codes)
if results:
log(f"Preparing to send {len(results)} records to DWH")
from etl_common import process_etls
process_etls(
application_name=APPLICATION_NAME,
etl_desc_url=f'{EROUTER_HOST}/procedure/{PRIMARY_MADAB}/umdap.etl.description/call.json?application={APPLICATION_NAME}',
data_type=None,
api_requestor=gtt_request,
erouter_url_list=UMDAP_ETL_URL_LIST
)
else:
log("No data to send")
except Exception as e:
log(f'Main thread unexpected error: {traceback.format_exception(*sys.exc_info())}')
sys.exit(1)
finally:
log('--------------------- Script finished -----------------------------')
if __name__ == "__main__":
main()