Refactor: Code-Qualität verbessert und Projektstruktur aufgeräumt
Some checks failed
Deployment / deploy-docker (push) Has been cancelled
Some checks failed
Deployment / deploy-docker (push) Has been cancelled
- daemon.py: gc.collect() entfernt, robustes Scheduling (last_run_date statt Minuten-Check), Exchange Registry Pattern eingeführt (STREAMING_EXCHANGES/STANDARD_EXCHANGES) - deutsche_boerse.py: Thread-safe User-Agent Rotation bei Rate-Limits, Logging statt print(), Feiertags-Prüfung, aufgeteilte Parse-Methoden - eix.py: Logging statt print(), spezifische Exception-Typen statt blankem except - read.py gelöscht und durch scripts/inspect_gzip.py ersetzt (Streaming-basiert) - Utility-Scripts in scripts/ verschoben (cleanup_duplicates, restore_and_fix, verify_fix)
This commit is contained in:
@@ -2,16 +2,20 @@ import requests
|
||||
import gzip
|
||||
import json
|
||||
import io
|
||||
import re
|
||||
import time
|
||||
import logging
|
||||
import threading
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import List, Optional
|
||||
from .base import BaseExchange, Trade
|
||||
from bs4 import BeautifulSoup
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Rate-Limiting Konfiguration
|
||||
RATE_LIMIT_DELAY = 0.5 # Sekunden zwischen Requests
|
||||
RATE_LIMIT_RETRY_DELAY = 5 # Sekunden Wartezeit bei 429
|
||||
MAX_RETRIES = 3 # Maximale Wiederholungen bei 429
|
||||
MAX_RETRIES = 5 # Maximale Wiederholungen bei 429
|
||||
|
||||
# API URLs für Deutsche Börse
|
||||
API_URLS = {
|
||||
@@ -21,17 +25,47 @@ API_URLS = {
|
||||
}
|
||||
DOWNLOAD_BASE_URL = "https://mfs.deutsche-boerse.com/api/download"
|
||||
|
||||
# Browser User-Agent für Zugriff
|
||||
HEADERS = {
|
||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Accept': 'application/json, application/gzip, */*',
|
||||
'Referer': 'https://mfs.deutsche-boerse.com/',
|
||||
}
|
||||
# Liste von User-Agents für Rotation bei Rate-Limiting
|
||||
USER_AGENTS = [
|
||||
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/121.0.0.0 Safari/537.36',
|
||||
'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:122.0) Gecko/20100101 Firefox/122.0',
|
||||
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.2 Safari/605.1.15',
|
||||
'Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36 Edg/119.0.0.0',
|
||||
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10.15; rv:121.0) Gecko/20100101 Firefox/121.0',
|
||||
'Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:122.0) Gecko/20100101 Firefox/122.0',
|
||||
]
|
||||
|
||||
|
||||
class UserAgentRotator:
|
||||
"""Thread-safe User-Agent Rotation"""
|
||||
|
||||
def __init__(self):
|
||||
self._index = 0
|
||||
self._lock = threading.Lock()
|
||||
|
||||
def get_headers(self, rotate: bool = False) -> dict:
|
||||
"""Gibt Headers mit aktuellem User-Agent zurück. Bei rotate=True wird zum nächsten gewechselt."""
|
||||
with self._lock:
|
||||
if rotate:
|
||||
self._index = (self._index + 1) % len(USER_AGENTS)
|
||||
return {
|
||||
'User-Agent': USER_AGENTS[self._index],
|
||||
'Accept': 'application/json, application/gzip, */*',
|
||||
'Referer': 'https://mfs.deutsche-boerse.com/',
|
||||
}
|
||||
|
||||
# Globale Instanz für User-Agent Rotation
|
||||
_ua_rotator = UserAgentRotator()
|
||||
|
||||
|
||||
class DeutscheBoerseBase(BaseExchange):
|
||||
"""Basisklasse für Deutsche Börse Exchanges (Xetra, Frankfurt, Quotrix)"""
|
||||
|
||||
# Regex für Dateinamen-Parsing (kompiliert für Performance)
|
||||
_FILENAME_PATTERN = re.compile(r'posttrade-(\d{4}-\d{2}-\d{2})T(\d{2})_(\d{2})')
|
||||
|
||||
@property
|
||||
def base_url(self) -> str:
|
||||
"""Override in subclasses"""
|
||||
@@ -46,60 +80,73 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
"""API URL für die Dateiliste"""
|
||||
return API_URLS.get(self.name, self.base_url)
|
||||
|
||||
def _handle_rate_limit(self, retry: int, context: str) -> None:
|
||||
"""Zentrale Rate-Limit Behandlung: rotiert User-Agent und wartet."""
|
||||
_ua_rotator.get_headers(rotate=True)
|
||||
wait_time = RATE_LIMIT_RETRY_DELAY * (retry + 1)
|
||||
logger.warning(f"[{self.name}] Rate limited ({context}), rotating User-Agent and waiting {wait_time}s... (retry {retry + 1}/{MAX_RETRIES})")
|
||||
time.sleep(wait_time)
|
||||
|
||||
def _get_file_list(self) -> List[str]:
|
||||
"""Holt die Dateiliste von der JSON API"""
|
||||
try:
|
||||
api_url = self.api_url
|
||||
print(f"[{self.name}] Fetching file list from: {api_url}")
|
||||
response = requests.get(api_url, headers=HEADERS, timeout=30)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
files = data.get('CurrentFiles', [])
|
||||
|
||||
print(f"[{self.name}] API returned {len(files)} files")
|
||||
if files:
|
||||
print(f"[{self.name}] Sample files: {files[:3]}")
|
||||
return files
|
||||
|
||||
except Exception as e:
|
||||
print(f"[{self.name}] Error fetching file list from API: {e}")
|
||||
import traceback
|
||||
print(f"[{self.name}] Traceback: {traceback.format_exc()}")
|
||||
return []
|
||||
api_url = self.api_url
|
||||
|
||||
for retry in range(MAX_RETRIES):
|
||||
try:
|
||||
headers = _ua_rotator.get_headers(rotate=(retry > 0))
|
||||
logger.info(f"[{self.name}] Fetching file list from: {api_url}")
|
||||
response = requests.get(api_url, headers=headers, timeout=30)
|
||||
|
||||
if response.status_code == 429:
|
||||
self._handle_rate_limit(retry, "file list")
|
||||
continue
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
files = data.get('CurrentFiles', [])
|
||||
|
||||
logger.info(f"[{self.name}] API returned {len(files)} files")
|
||||
if files:
|
||||
logger.debug(f"[{self.name}] Sample files: {files[:3]}")
|
||||
return files
|
||||
|
||||
except requests.exceptions.HTTPError as e:
|
||||
if e.response.status_code == 429:
|
||||
self._handle_rate_limit(retry, "file list HTTPError")
|
||||
continue
|
||||
logger.error(f"[{self.name}] HTTP error fetching file list: {e}")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.exception(f"[{self.name}] Error fetching file list from API: {e}")
|
||||
break
|
||||
|
||||
return []
|
||||
|
||||
def _filter_files_for_date(self, files: List[str], target_date: datetime.date) -> List[str]:
|
||||
"""
|
||||
Filtert Dateien für ein bestimmtes Datum.
|
||||
Dateiformat: DETR-posttrade-YYYY-MM-DDTHH_MM.json.gz (mit Unterstrich!)
|
||||
Dateiformat: DETR-posttrade-YYYY-MM-DDTHH_MM.json.gz
|
||||
|
||||
Da Handel bis 22:00 MEZ geht (21:00/20:00 UTC), müssen wir auch
|
||||
Dateien nach Mitternacht UTC berücksichtigen.
|
||||
"""
|
||||
import re
|
||||
filtered = []
|
||||
|
||||
# Für den Vortag: Dateien vom target_date UND vom Folgetag (bis ~02:00 UTC)
|
||||
target_str = target_date.strftime('%Y-%m-%d')
|
||||
next_day = target_date + timedelta(days=1)
|
||||
next_day_str = next_day.strftime('%Y-%m-%d')
|
||||
|
||||
for file in files:
|
||||
# Extrahiere Datum aus Dateiname
|
||||
# Format: DETR-posttrade-2026-01-26T21_30.json.gz
|
||||
if target_str in file:
|
||||
filtered.append(file)
|
||||
elif next_day_str in file:
|
||||
# Prüfe ob es eine frühe Datei vom nächsten Tag ist (< 03:00 UTC)
|
||||
try:
|
||||
# Finde Timestamp im Dateinamen mit Unterstrich für Minuten
|
||||
match = re.search(r'posttrade-(\d{4}-\d{2}-\d{2})T(\d{2})_(\d{2})', file)
|
||||
if match:
|
||||
hour = int(match.group(2))
|
||||
if hour < 3: # Frühe Morgenstunden gehören noch zum Vortag
|
||||
filtered.append(file)
|
||||
except Exception:
|
||||
pass
|
||||
match = self._FILENAME_PATTERN.search(file)
|
||||
if match:
|
||||
hour = int(match.group(2))
|
||||
if hour < 3: # Frühe Morgenstunden gehören noch zum Vortag
|
||||
filtered.append(file)
|
||||
|
||||
return filtered
|
||||
|
||||
@@ -110,17 +157,14 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
|
||||
for retry in range(MAX_RETRIES):
|
||||
try:
|
||||
response = requests.get(full_url, headers=HEADERS, timeout=60)
|
||||
headers = _ua_rotator.get_headers(rotate=(retry > 0))
|
||||
response = requests.get(full_url, headers=headers, timeout=60)
|
||||
|
||||
if response.status_code == 404:
|
||||
# Datei nicht gefunden - normal für alte Dateien
|
||||
return []
|
||||
|
||||
if response.status_code == 429:
|
||||
# Rate-Limit erreicht - warten und erneut versuchen
|
||||
wait_time = RATE_LIMIT_RETRY_DELAY * (retry + 1)
|
||||
print(f"[{self.name}] Rate limited, waiting {wait_time}s...")
|
||||
time.sleep(wait_time)
|
||||
self._handle_rate_limit(retry, "download")
|
||||
continue
|
||||
|
||||
response.raise_for_status()
|
||||
@@ -130,13 +174,11 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
content = f.read().decode('utf-8')
|
||||
|
||||
if not content.strip():
|
||||
# Leere Datei
|
||||
return []
|
||||
|
||||
# NDJSON Format: Eine JSON-Zeile pro Trade
|
||||
lines = content.strip().split('\n')
|
||||
if not lines or (len(lines) == 1 and not lines[0].strip()):
|
||||
# Leere Datei
|
||||
return []
|
||||
|
||||
for line in lines:
|
||||
@@ -147,116 +189,146 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
trade = self._parse_trade_record(record)
|
||||
if trade:
|
||||
trades.append(trade)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
except Exception:
|
||||
continue
|
||||
except json.JSONDecodeError as e:
|
||||
logger.debug(f"[{self.name}] JSON decode error in {filename}: {e}")
|
||||
except Exception as e:
|
||||
logger.debug(f"[{self.name}] Error parsing record in {filename}: {e}")
|
||||
|
||||
# Erfolg - keine weitere Retry nötig
|
||||
# Erfolg
|
||||
break
|
||||
|
||||
except requests.exceptions.HTTPError as e:
|
||||
if e.response.status_code == 429:
|
||||
wait_time = RATE_LIMIT_RETRY_DELAY * (retry + 1)
|
||||
print(f"[{self.name}] Rate limited, waiting {wait_time}s...")
|
||||
time.sleep(wait_time)
|
||||
self._handle_rate_limit(retry, "download HTTPError")
|
||||
continue
|
||||
elif e.response.status_code != 404:
|
||||
print(f"[{self.name}] HTTP error downloading {filename}: {e}")
|
||||
logger.error(f"[{self.name}] HTTP error downloading {filename}: {e}")
|
||||
break
|
||||
except Exception as e:
|
||||
print(f"[{self.name}] Error downloading/parsing {filename}: {e}")
|
||||
logger.error(f"[{self.name}] Error downloading/parsing {filename}: {e}")
|
||||
break
|
||||
|
||||
return trades
|
||||
|
||||
def _parse_timestamp(self, ts_str: str) -> Optional[datetime]:
|
||||
"""
|
||||
Parst einen Timestamp-String in ein datetime-Objekt.
|
||||
Unterstützt Nanosekunden durch Kürzung auf Mikrosekunden.
|
||||
"""
|
||||
if not ts_str:
|
||||
return None
|
||||
|
||||
# Ersetze 'Z' durch '+00:00' für ISO-Kompatibilität
|
||||
ts_str = ts_str.replace('Z', '+00:00')
|
||||
|
||||
# Kürze Nanosekunden auf Mikrosekunden (Python max 6 Dezimalstellen)
|
||||
if '.' in ts_str:
|
||||
# Split bei '+' oder '-' für Timezone
|
||||
if '+' in ts_str:
|
||||
time_part, tz_part = ts_str.rsplit('+', 1)
|
||||
tz_part = '+' + tz_part
|
||||
elif ts_str.count('-') > 2: # Negative Timezone
|
||||
time_part, tz_part = ts_str.rsplit('-', 1)
|
||||
tz_part = '-' + tz_part
|
||||
else:
|
||||
time_part, tz_part = ts_str, ''
|
||||
|
||||
if '.' in time_part:
|
||||
base, frac = time_part.split('.')
|
||||
frac = frac[:6] # Kürze auf 6 Stellen
|
||||
ts_str = f"{base}.{frac}{tz_part}"
|
||||
|
||||
return datetime.fromisoformat(ts_str)
|
||||
|
||||
def _extract_price(self, record: dict) -> Optional[float]:
|
||||
"""Extrahiert den Preis aus verschiedenen JSON-Formaten."""
|
||||
# Neues Format
|
||||
if 'lastTrade' in record:
|
||||
return float(record['lastTrade'])
|
||||
|
||||
# Altes Format mit verschachteltem Pric-Objekt
|
||||
pric = record.get('Pric')
|
||||
if pric is None:
|
||||
return None
|
||||
|
||||
if isinstance(pric, (int, float)):
|
||||
return float(pric)
|
||||
|
||||
if isinstance(pric, dict):
|
||||
# Versuche verschiedene Pfade
|
||||
if 'Pric' in pric:
|
||||
inner = pric['Pric']
|
||||
if isinstance(inner, dict):
|
||||
amt = inner.get('MntryVal', {}).get('Amt') or inner.get('Amt')
|
||||
if amt is not None:
|
||||
return float(amt)
|
||||
if 'MntryVal' in pric:
|
||||
amt = pric['MntryVal'].get('Amt')
|
||||
if amt is not None:
|
||||
return float(amt)
|
||||
|
||||
return None
|
||||
|
||||
def _extract_quantity(self, record: dict) -> Optional[float]:
|
||||
"""Extrahiert die Menge aus verschiedenen JSON-Formaten."""
|
||||
# Neues Format
|
||||
if 'lastQty' in record:
|
||||
return float(record['lastQty'])
|
||||
|
||||
# Altes Format
|
||||
qty = record.get('Qty')
|
||||
if qty is None:
|
||||
return None
|
||||
|
||||
if isinstance(qty, (int, float)):
|
||||
return float(qty)
|
||||
|
||||
if isinstance(qty, dict):
|
||||
val = qty.get('Unit') or qty.get('Qty')
|
||||
if val is not None:
|
||||
return float(val)
|
||||
|
||||
return None
|
||||
|
||||
def _parse_trade_record(self, record: dict) -> Optional[Trade]:
|
||||
"""
|
||||
Parst einen einzelnen Trade-Record aus dem JSON.
|
||||
|
||||
Aktuelles JSON-Format (NDJSON):
|
||||
{
|
||||
"messageId": "posttrade",
|
||||
"sourceName": "GAT",
|
||||
"isin": "US00123Q1040",
|
||||
"lastTradeTime": "2026-01-29T14:07:00.419000000Z",
|
||||
"lastTrade": 10.145,
|
||||
"lastQty": 500.0,
|
||||
"currency": "EUR",
|
||||
...
|
||||
}
|
||||
Unterstützte Formate:
|
||||
- Neues Format: isin, lastTrade, lastQty, lastTradeTime
|
||||
- Altes Format: FinInstrmId.Id, Pric, Qty, TrdDt/TrdTm
|
||||
"""
|
||||
try:
|
||||
# ISIN extrahieren - neues Format verwendet 'isin' lowercase
|
||||
isin = record.get('isin') or record.get('ISIN') or record.get('instrumentId') or record.get('FinInstrmId', {}).get('Id', '')
|
||||
# ISIN extrahieren
|
||||
isin = (
|
||||
record.get('isin') or
|
||||
record.get('ISIN') or
|
||||
record.get('instrumentId') or
|
||||
record.get('FinInstrmId', {}).get('Id', '')
|
||||
)
|
||||
if not isin:
|
||||
return None
|
||||
|
||||
# Preis extrahieren - neues Format: 'lastTrade'
|
||||
price = None
|
||||
if 'lastTrade' in record:
|
||||
price = float(record['lastTrade'])
|
||||
elif 'Pric' in record:
|
||||
pric = record['Pric']
|
||||
if isinstance(pric, dict):
|
||||
if 'Pric' in pric:
|
||||
inner = pric['Pric']
|
||||
if 'MntryVal' in inner:
|
||||
price = float(inner['MntryVal'].get('Amt', 0))
|
||||
elif 'Amt' in inner:
|
||||
price = float(inner['Amt'])
|
||||
elif 'MntryVal' in pric:
|
||||
price = float(pric['MntryVal'].get('Amt', 0))
|
||||
elif isinstance(pric, (int, float)):
|
||||
price = float(pric)
|
||||
|
||||
# Preis extrahieren
|
||||
price = self._extract_price(record)
|
||||
if price is None or price <= 0:
|
||||
return None
|
||||
|
||||
# Menge extrahieren - neues Format: 'lastQty'
|
||||
quantity = None
|
||||
if 'lastQty' in record:
|
||||
quantity = float(record['lastQty'])
|
||||
elif 'Qty' in record:
|
||||
qty = record['Qty']
|
||||
if isinstance(qty, dict):
|
||||
quantity = float(qty.get('Unit', qty.get('Qty', 0)))
|
||||
elif isinstance(qty, (int, float)):
|
||||
quantity = float(qty)
|
||||
|
||||
# Menge extrahieren
|
||||
quantity = self._extract_quantity(record)
|
||||
if quantity is None or quantity <= 0:
|
||||
return None
|
||||
|
||||
# Timestamp extrahieren - neues Format: 'lastTradeTime'
|
||||
# Timestamp extrahieren
|
||||
timestamp = None
|
||||
if 'lastTradeTime' in record:
|
||||
ts_str = record['lastTradeTime']
|
||||
# Format: "2026-01-29T14:07:00.419000000Z"
|
||||
# Python kann max 6 Dezimalstellen, also kürzen
|
||||
if '.' in ts_str:
|
||||
parts = ts_str.replace('Z', '').split('.')
|
||||
if len(parts) == 2 and len(parts[1]) > 6:
|
||||
ts_str = parts[0] + '.' + parts[1][:6] + '+00:00'
|
||||
else:
|
||||
ts_str = ts_str.replace('Z', '+00:00')
|
||||
else:
|
||||
ts_str = ts_str.replace('Z', '+00:00')
|
||||
timestamp = datetime.fromisoformat(ts_str)
|
||||
timestamp = self._parse_timestamp(record['lastTradeTime'])
|
||||
else:
|
||||
# Fallback für altes Format
|
||||
trd_dt = record.get('TrdDt', '')
|
||||
trd_tm = record.get('TrdTm', '00:00:00')
|
||||
|
||||
if not trd_dt:
|
||||
return None
|
||||
|
||||
ts_str = f"{trd_dt}T{trd_tm}"
|
||||
if '.' in ts_str:
|
||||
parts = ts_str.split('.')
|
||||
if len(parts[1]) > 6:
|
||||
ts_str = parts[0] + '.' + parts[1][:6]
|
||||
|
||||
timestamp = datetime.fromisoformat(ts_str)
|
||||
if trd_dt:
|
||||
timestamp = self._parse_timestamp(f"{trd_dt}T{trd_tm}")
|
||||
|
||||
if timestamp is None:
|
||||
return None
|
||||
@@ -273,22 +345,41 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
timestamp=timestamp
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
# Debug: Zeige ersten fehlgeschlagenen Record
|
||||
except (ValueError, TypeError, KeyError) as e:
|
||||
logger.debug(f"[{self.name}] Failed to parse trade record: {e}")
|
||||
return None
|
||||
|
||||
def _get_last_trading_day(self, from_date: datetime.date) -> datetime.date:
|
||||
"""
|
||||
Findet den letzten Handelstag (überspringt Wochenenden).
|
||||
Findet den letzten Handelstag (überspringt Wochenenden und bekannte Feiertage).
|
||||
Montag=0, Sonntag=6
|
||||
"""
|
||||
# Deutsche Börsen-Feiertage (fixe Daten, jedes Jahr gleich)
|
||||
# Bewegliche Feiertage (Ostern etc.) müssten jährlich berechnet werden
|
||||
fixed_holidays = {
|
||||
(1, 1), # Neujahr
|
||||
(5, 1), # Tag der Arbeit
|
||||
(12, 24), # Heiligabend
|
||||
(12, 25), # 1. Weihnachtstag
|
||||
(12, 26), # 2. Weihnachtstag
|
||||
(12, 31), # Silvester
|
||||
}
|
||||
|
||||
date = from_date
|
||||
# Wenn Samstag (5), gehe zurück zu Freitag
|
||||
if date.weekday() == 5:
|
||||
date = date - timedelta(days=1)
|
||||
# Wenn Sonntag (6), gehe zurück zu Freitag
|
||||
elif date.weekday() == 6:
|
||||
date = date - timedelta(days=2)
|
||||
max_iterations = 10 # Sicherheit gegen Endlosschleife
|
||||
|
||||
for _ in range(max_iterations):
|
||||
# Wochenende überspringen
|
||||
if date.weekday() == 5: # Samstag
|
||||
date = date - timedelta(days=1)
|
||||
elif date.weekday() == 6: # Sonntag
|
||||
date = date - timedelta(days=2)
|
||||
# Feiertag überspringen
|
||||
elif (date.month, date.day) in fixed_holidays:
|
||||
date = date - timedelta(days=1)
|
||||
else:
|
||||
break
|
||||
|
||||
return date
|
||||
|
||||
def fetch_latest_trades(self, include_yesterday: bool = True, since_date: datetime = None) -> List[Trade]:
|
||||
@@ -304,40 +395,36 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
# Standard: Vortag
|
||||
target_date = (datetime.now(timezone.utc) - timedelta(days=1)).date()
|
||||
|
||||
# Überspringe Wochenenden
|
||||
# Überspringe Wochenenden und Feiertage
|
||||
original_date = target_date
|
||||
target_date = self._get_last_trading_day(target_date)
|
||||
|
||||
if target_date != original_date:
|
||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||
logger.info(f"[{self.name}] Adjusted date: {original_date} -> {target_date} (weekend/holiday)")
|
||||
|
||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||
|
||||
# Hole Dateiliste von der API
|
||||
files = self._get_file_list()
|
||||
|
||||
if not files:
|
||||
print(f"[{self.name}] No files available from API")
|
||||
logger.warning(f"[{self.name}] No files available from API")
|
||||
return []
|
||||
|
||||
# Dateien für Zieldatum filtern
|
||||
target_files = self._filter_files_for_date(files, target_date)
|
||||
print(f"[{self.name}] {len(target_files)} files match target date (of {len(files)} total)")
|
||||
logger.info(f"[{self.name}] {len(target_files)} files match target date (of {len(files)} total)")
|
||||
|
||||
if not target_files:
|
||||
print(f"[{self.name}] No files for target date found")
|
||||
logger.warning(f"[{self.name}] No files for target date found")
|
||||
return []
|
||||
|
||||
# Alle passenden Dateien herunterladen und parsen (mit Rate-Limiting)
|
||||
# Alle passenden Dateien herunterladen und parsen
|
||||
successful = 0
|
||||
failed = 0
|
||||
total_files = len(target_files)
|
||||
|
||||
if total_files == 0:
|
||||
print(f"[{self.name}] No files to download for date {target_date}")
|
||||
return []
|
||||
|
||||
print(f"[{self.name}] Starting download of {total_files} files...")
|
||||
logger.info(f"[{self.name}] Starting download of {total_files} files...")
|
||||
|
||||
for i, file in enumerate(target_files):
|
||||
trades = self._download_and_parse_file(file)
|
||||
@@ -353,9 +440,9 @@ class DeutscheBoerseBase(BaseExchange):
|
||||
|
||||
# Fortschritt alle 100 Dateien
|
||||
if (i + 1) % 100 == 0:
|
||||
print(f"[{self.name}] Progress: {i + 1}/{total_files} files, {successful} successful, {len(all_trades)} trades so far")
|
||||
logger.info(f"[{self.name}] Progress: {i + 1}/{total_files} files, {successful} successful, {len(all_trades)} trades so far")
|
||||
|
||||
print(f"[{self.name}] Downloaded {successful} files ({failed} failed/empty), total {len(all_trades)} trades")
|
||||
logger.info(f"[{self.name}] Downloaded {successful} files ({failed} failed/empty), total {len(all_trades)} trades")
|
||||
return all_trades
|
||||
|
||||
|
||||
|
||||
@@ -1,29 +1,38 @@
|
||||
import requests
|
||||
import json
|
||||
from bs4 import BeautifulSoup
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
from typing import List, Generator, Tuple, Optional
|
||||
from typing import List, Generator, Tuple
|
||||
from .base import BaseExchange, Trade
|
||||
import csv
|
||||
import io
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class EIXExchange(BaseExchange):
|
||||
"""European Investor Exchange - CSV-basierte Trade-Daten."""
|
||||
|
||||
API_BASE_URL = "https://european-investor-exchange.com/api"
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
return "EIX"
|
||||
|
||||
def get_files_to_process(self, limit: int = 1, since_date: datetime = None) -> List[dict]:
|
||||
"""Holt die Liste der zu verarbeitenden Dateien ohne sie herunterzuladen."""
|
||||
url = "https://european-investor-exchange.com/api/official-trades"
|
||||
url = f"{self.API_BASE_URL}/official-trades"
|
||||
try:
|
||||
response = requests.get(url, timeout=15)
|
||||
response.raise_for_status()
|
||||
files_list = response.json()
|
||||
except Exception as e:
|
||||
print(f"Error fetching EIX file list: {e}")
|
||||
except requests.exceptions.RequestException as e:
|
||||
logger.error(f"[{self.name}] Fehler beim Abrufen der Dateiliste: {e}")
|
||||
return []
|
||||
except ValueError as e:
|
||||
logger.error(f"[{self.name}] Ungültiges JSON in Dateiliste: {e}")
|
||||
return []
|
||||
|
||||
# Filter files based on date in filename if since_date provided
|
||||
# Filtere Dateien nach Datum wenn since_date angegeben
|
||||
filtered_files = []
|
||||
for item in files_list:
|
||||
file_key = item.get('fileName')
|
||||
@@ -39,7 +48,9 @@ class EIXExchange(BaseExchange):
|
||||
|
||||
if file_date.date() >= since_date.date():
|
||||
filtered_files.append(item)
|
||||
except Exception:
|
||||
except (ValueError, IndexError) as e:
|
||||
# Dateiname hat unerwartetes Format - zur Sicherheit einschließen
|
||||
logger.debug(f"[{self.name}] Konnte Datum nicht aus {file_key} extrahieren: {e}")
|
||||
filtered_files.append(item)
|
||||
else:
|
||||
filtered_files.append(item)
|
||||
@@ -57,13 +68,15 @@ class EIXExchange(BaseExchange):
|
||||
if not file_key:
|
||||
return []
|
||||
|
||||
csv_url = f"https://european-investor-exchange.com/api/trade-file-contents?key={file_key}"
|
||||
csv_url = f"{self.API_BASE_URL}/trade-file-contents?key={file_key}"
|
||||
try:
|
||||
csv_response = requests.get(csv_url, timeout=60)
|
||||
if csv_response.status_code == 200:
|
||||
return self._parse_csv(csv_response.text)
|
||||
response = requests.get(csv_url, timeout=60)
|
||||
response.raise_for_status()
|
||||
return self._parse_csv(response.text)
|
||||
except requests.exceptions.RequestException as e:
|
||||
logger.error(f"[{self.name}] Fehler beim Download von {file_key}: {e}")
|
||||
except Exception as e:
|
||||
print(f"Error downloading EIX CSV {file_key}: {e}")
|
||||
logger.error(f"[{self.name}] Unerwarteter Fehler bei {file_key}: {e}")
|
||||
|
||||
return []
|
||||
|
||||
@@ -80,7 +93,7 @@ class EIXExchange(BaseExchange):
|
||||
if trades:
|
||||
yield (file_key, trades)
|
||||
|
||||
def fetch_latest_trades(self, limit: int = 1, since_date: datetime = None) -> List[Trade]:
|
||||
def fetch_latest_trades(self, limit: int = 1, since_date: datetime = None, **kwargs) -> List[Trade]:
|
||||
"""
|
||||
Legacy-Methode für Kompatibilität.
|
||||
WARNUNG: Lädt alle Trades in den Speicher! Für große Datenmengen fetch_trades_streaming() verwenden.
|
||||
@@ -93,33 +106,53 @@ class EIXExchange(BaseExchange):
|
||||
return all_trades
|
||||
|
||||
# Für große Requests: Warnung ausgeben und leere Liste zurückgeben
|
||||
# Der Daemon soll stattdessen fetch_trades_streaming() verwenden
|
||||
print(f"[EIX] WARNING: fetch_latest_trades() called with large dataset. Use streaming instead.")
|
||||
logger.warning(f"[{self.name}] fetch_latest_trades() mit großem Dataset aufgerufen. Verwende Streaming.")
|
||||
return []
|
||||
|
||||
def _parse_csv(self, csv_text: str) -> List[Trade]:
|
||||
"""Parst CSV-Text zu Trade-Objekten."""
|
||||
trades = []
|
||||
parse_errors = 0
|
||||
|
||||
f = io.StringIO(csv_text)
|
||||
reader = csv.DictReader(f, delimiter=',')
|
||||
for row in reader:
|
||||
|
||||
for row_num, row in enumerate(reader, start=2): # Start bei 2 wegen Header
|
||||
try:
|
||||
price = float(row['Unit Price'])
|
||||
quantity = float(row['Quantity'])
|
||||
isin = row['Instrument Identifier']
|
||||
symbol = isin
|
||||
time_str = row['Trading day & Trading time UTC']
|
||||
|
||||
# Preis und Menge validieren
|
||||
if price <= 0 or quantity <= 0:
|
||||
logger.debug(f"[{self.name}] Zeile {row_num}: Ungültiger Preis/Menge: {price}/{quantity}")
|
||||
parse_errors += 1
|
||||
continue
|
||||
|
||||
ts_str = time_str.replace('Z', '+00:00')
|
||||
timestamp = datetime.fromisoformat(ts_str)
|
||||
|
||||
trades.append(Trade(
|
||||
exchange=self.name,
|
||||
symbol=symbol,
|
||||
symbol=isin,
|
||||
isin=isin,
|
||||
price=price,
|
||||
quantity=quantity,
|
||||
timestamp=timestamp
|
||||
))
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
except KeyError as e:
|
||||
logger.debug(f"[{self.name}] Zeile {row_num}: Fehlendes Feld {e}")
|
||||
parse_errors += 1
|
||||
except ValueError as e:
|
||||
logger.debug(f"[{self.name}] Zeile {row_num}: Ungültiger Wert: {e}")
|
||||
parse_errors += 1
|
||||
except Exception as e:
|
||||
logger.warning(f"[{self.name}] Zeile {row_num}: Unerwarteter Fehler: {e}")
|
||||
parse_errors += 1
|
||||
|
||||
if parse_errors > 0:
|
||||
logger.debug(f"[{self.name}] {parse_errors} Zeilen konnten nicht geparst werden")
|
||||
|
||||
return trades
|
||||
|
||||
Reference in New Issue
Block a user