Compare commits
9 Commits
f0acbf39bc
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ad890d23b0 | ||
|
|
819b6c2232 | ||
|
|
24da2a5861 | ||
|
|
18801ec27e | ||
|
|
8634c01ec0 | ||
|
|
71f8614cb5 | ||
|
|
80d8801728 | ||
|
|
846f5e76fe | ||
|
|
7c802a3907 |
@@ -7,4 +7,6 @@ RUN pip install --no-cache-dir -r requirements.txt
|
|||||||
|
|
||||||
COPY . .
|
COPY . .
|
||||||
|
|
||||||
|
ENV PYTHONPATH=/app
|
||||||
|
|
||||||
CMD ["python", "dashboard/server.py"]
|
CMD ["python", "dashboard/server.py"]
|
||||||
|
|||||||
171
README.md
171
README.md
@@ -1,20 +1,140 @@
|
|||||||
# Trading Data Daemon
|
# Trading Data Daemon
|
||||||
|
|
||||||
Ein modularer Daemon zum Herunterladen und Speichern von Handelsdaten von verschiedenen Börsen in einer Time-Series-Datenbank.
|
Ein modularer Daemon zum Herunterladen und Speichern von Handelsdaten von verschiedenen deutschen Börsen in QuestDB (Time-Series-Datenbank).
|
||||||
|
|
||||||
## Unterstützte Exchanges
|
## Unterstützte Exchanges
|
||||||
- **European Investor Exchange (EIX)**: Lädt tägliche Kursblatt-CSVs herunter.
|
|
||||||
- **Lang & Schwarz (LS)**: Fragt die heutigen Trades über deren JSON/CSV-RPC ab.
|
- **European Investor Exchange (EIX)** — Streaming-Verarbeitung (große CSV-Dateien)
|
||||||
|
- **Lang & Schwarz (LS)** — JSON/CSV-RPC API
|
||||||
|
- **Deutsche Börse** — Xetra, Frankfurt, Quotrix
|
||||||
|
- **Gettex** — Bayerische Börse (MUNC/MUND)
|
||||||
|
- **Stuttgart** — MiFIR II Delayed Data
|
||||||
|
- **Börsenag** — Düsseldorf (DUSA/DUSB/DUSC/DUSD), Hamburg (HAMA/HAMB), Hannover (HANA/HANB)
|
||||||
|
|
||||||
## Architektur
|
## Architektur
|
||||||
- `src/exchanges/base.py`: Basisklasse für neue Börsen (einfach erweiterbar).
|
|
||||||
- `src/database/questdb_client.py`: Speichert Daten in QuestDB via Influx Line Protocol (ILP).
|
Fünf Microservices, orchestriert via Docker Compose:
|
||||||
- `daemon.py`: Der Orchestrator, der die Daten abruft und speichert.
|
|
||||||
|
```
|
||||||
|
┌─────────────────────────────────────────────────────────────────┐
|
||||||
|
│ QuestDB │
|
||||||
|
│ (9000=HTTP, 8812=PostgreSQL, 9009=ILP) │
|
||||||
|
└──────────┬──────────────┬──────────────┬──────────────┬─────────┘
|
||||||
|
│ │ │ │
|
||||||
|
┌──────┴──────┐ ┌─────┴─────┐ ┌──────┴──────┐ ┌─────┴─────┐
|
||||||
|
│ fetcher │ │ analytics │ │ metadata │ │ dashboard │
|
||||||
|
│ daemon.py │ │ worker │ │ fetcher │ │ :8080 │
|
||||||
|
└─────────────┘ └───────────┘ └─────────────┘ └───────────┘
|
||||||
|
```
|
||||||
|
|
||||||
|
- **Fetcher (`daemon.py`):** Orchestrator. Holt Trades von allen Börsen täglich um 23:00. Streaming für EIX, Batch für alle anderen. Deduplizierung via Hash-Cache.
|
||||||
|
- **Analytics Worker (`src/analytics/worker.py`):** Berechnet aggregierte Tabellen für Zeiträume: 7, 30, 42, 69, 180, 365 Tage.
|
||||||
|
- **Metadata Fetcher (`src/metadata/fetcher.py`):** Reichert ISINs mit Firmen-/Sektordaten an (OpenFIGI API, yfinance).
|
||||||
|
- **Dashboard (`dashboard/server.py`):** FastAPI-Server mit REST-Endpunkten und statischem UI aus `dashboard/public/`.
|
||||||
|
|
||||||
|
## Datenbank-Schema
|
||||||
|
|
||||||
|
QuestDB speichert alle Daten via Influx Line Protocol (ILP) mit Nanosekunden-Präzision.
|
||||||
|
|
||||||
|
### `trades`
|
||||||
|
|
||||||
|
Rohe Handelsdaten aller Börsen. Geschrieben vom Fetcher.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Zeitpunkt des Trades |
|
||||||
|
| `exchange` | symbol (tag) | Börsenname (z.B. `XETRA`, `LS`, `GETTEX`) |
|
||||||
|
| `symbol` | symbol (tag) | Wertpapiername |
|
||||||
|
| `isin` | symbol (tag) | ISIN-Kennung (z.B. `DE000BAY0017`) |
|
||||||
|
| `price` | double | Handelspreis |
|
||||||
|
| `quantity` | double | Handelsvolumen (Stückzahl) |
|
||||||
|
|
||||||
|
### `analytics_exchange_daily`
|
||||||
|
|
||||||
|
Tägliche Aggregationen pro Börse mit Moving Averages. Geschrieben vom Analytics Worker.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Tag der Aggregation |
|
||||||
|
| `exchange` | symbol (tag) | Börsenname |
|
||||||
|
| `trade_count` | long | Anzahl Trades am Tag |
|
||||||
|
| `volume` | double | Gesamtvolumen (Summe von Preis × Menge) |
|
||||||
|
| `ma{N}_count` | double | N-Tage Moving Average der Trade-Anzahl |
|
||||||
|
| `ma{N}_volume` | double | N-Tage Moving Average des Volumens |
|
||||||
|
|
||||||
|
*N = 7, 30, 42, 69, 180, 365*
|
||||||
|
|
||||||
|
### `analytics_daily_summary`
|
||||||
|
|
||||||
|
Tagesübergreifende Zusammenfassung aller Börsen. Geschrieben vom Analytics Worker.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Tag der Zusammenfassung |
|
||||||
|
| `total_trades` | long | Gesamtanzahl Trades über alle Börsen |
|
||||||
|
| `total_volume` | double | Gesamtvolumen über alle Börsen |
|
||||||
|
| `unique_assets` | long | Anzahl verschiedener gehandelter ISINs |
|
||||||
|
|
||||||
|
### `analytics_stock_trends`
|
||||||
|
|
||||||
|
Trendanalyse pro ISIN mit prozentualen Veränderungen. Geschrieben vom Analytics Worker.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Analysedatum |
|
||||||
|
| `isin` | symbol (tag) | ISIN-Kennung |
|
||||||
|
| `trade_count` | long | Gesamtanzahl Trades im Zeitraum |
|
||||||
|
| `volume` | double | Gesamtvolumen im Zeitraum |
|
||||||
|
| `count_change_pct` | double | Prozentuale Änderung der Trade-Anzahl (1. vs. 2. Hälfte) |
|
||||||
|
| `volume_change_pct` | double | Prozentuale Änderung des Volumens |
|
||||||
|
| `period_days` | long | Zeitraum in Tagen (7/30/42/69/180/365) |
|
||||||
|
|
||||||
|
### `analytics_volume_changes`
|
||||||
|
|
||||||
|
Volumen- und Anzahländerungen pro Börse mit Trendklassifizierung. Geschrieben vom Analytics Worker.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Analysedatum |
|
||||||
|
| `exchange` | symbol (tag) | Börsenname |
|
||||||
|
| `trend` | symbol (tag) | Trendklasse (s.u.) |
|
||||||
|
| `trade_count` | long | Gesamtanzahl Trades im Zeitraum |
|
||||||
|
| `volume` | double | Gesamtvolumen im Zeitraum |
|
||||||
|
| `count_change_pct` | double | Prozentuale Änderung der Trade-Anzahl |
|
||||||
|
| `volume_change_pct` | double | Prozentuale Änderung des Volumens |
|
||||||
|
| `period_days` | long | Zeitraum in Tagen |
|
||||||
|
|
||||||
|
**Trendklassen:** `mehr_trades_mehr_volumen`, `mehr_trades_weniger_volumen`, `weniger_trades_mehr_volumen`, `weniger_trades_weniger_volumen`, `stabil`
|
||||||
|
|
||||||
|
### `analytics_custom`
|
||||||
|
|
||||||
|
Vorberechnete Custom-Analytics für den Dashboard-Graphen-Builder. Geschrieben vom Analytics Worker.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Tag der Berechnung |
|
||||||
|
| `y_axis` | symbol (tag) | Metrik-Typ (`volume`, `trade_count`, `avg_price`) |
|
||||||
|
| `group_by` | symbol (tag) | Gruppierung (`exchange`, `isin`, `date`) |
|
||||||
|
| `exchange_filter` | symbol (tag) | Exchange-Filter (Börsenname oder `all`) |
|
||||||
|
| `group_value` | string | Wert der Gruppierung (z.B. Börsenname, ISIN) |
|
||||||
|
| `y_value` | double | Berechneter Metrik-Wert |
|
||||||
|
|
||||||
|
### `metadata`
|
||||||
|
|
||||||
|
ISIN-Stammdaten (Firmenname, Land, Sektor). Geschrieben vom Metadata Fetcher.
|
||||||
|
|
||||||
|
| Spalte | Typ | Beschreibung |
|
||||||
|
|--------|-----|-------------|
|
||||||
|
| `timestamp` | timestamp | Zeitpunkt der letzten Aktualisierung |
|
||||||
|
| `isin` | symbol (tag) | ISIN-Kennung |
|
||||||
|
| `name` | string | Firmen-/Wertpapiername |
|
||||||
|
| `country` | string | Ländercode oder -name |
|
||||||
|
| `continent` | string | Kontinent |
|
||||||
|
| `sector` | string | Branchenklassifizierung |
|
||||||
|
|
||||||
## Installation und Setup
|
## Installation und Setup
|
||||||
|
|
||||||
### 1. QuestDB (Timeseries DB) starten
|
### 1. QuestDB starten (via Docker Compose)
|
||||||
Am einfachsten via Docker Compose:
|
|
||||||
```bash
|
```bash
|
||||||
docker-compose up -d
|
docker-compose up -d
|
||||||
```
|
```
|
||||||
@@ -25,26 +145,35 @@ QuestDB ist dann unter `http://localhost:9000` erreichbar.
|
|||||||
pip install -r requirements.txt
|
pip install -r requirements.txt
|
||||||
```
|
```
|
||||||
|
|
||||||
### 3. Systemd Service einrichten
|
### 3. Manuell starten
|
||||||
Kopiere die Dateien nach `/etc/systemd/system/`:
|
```bash
|
||||||
|
python3 daemon.py # Fetcher
|
||||||
|
python -m src.analytics.worker # Analytics Worker
|
||||||
|
python src/metadata/fetcher.py # Metadata Fetcher
|
||||||
|
python dashboard/server.py # Dashboard (Port 8000)
|
||||||
|
```
|
||||||
|
|
||||||
|
### 4. Systemd Service (optional)
|
||||||
```bash
|
```bash
|
||||||
sudo cp systemd/trading-daemon.service /etc/systemd/system/
|
sudo cp systemd/trading-daemon.service /etc/systemd/system/
|
||||||
sudo cp systemd/trading-daemon.timer /etc/systemd/system/
|
sudo cp systemd/trading-daemon.timer /etc/systemd/system/
|
||||||
```
|
|
||||||
|
|
||||||
Pfade in `trading-daemon.service` müssen ggf. angepasst werden (aktuell auf `/Users/melchiorreimers/...` gesetzt).
|
|
||||||
|
|
||||||
Dienste aktivieren:
|
|
||||||
```bash
|
|
||||||
sudo systemctl daemon-reload
|
sudo systemctl daemon-reload
|
||||||
sudo systemctl enable --now trading-daemon.timer
|
sudo systemctl enable --now trading-daemon.timer
|
||||||
```
|
```
|
||||||
|
|
||||||
### 4. Manuell testen
|
## Umgebungsvariablen
|
||||||
```bash
|
|
||||||
python3 daemon.py
|
| Variable | Default | Beschreibung |
|
||||||
```
|
|----------|---------|-------------|
|
||||||
|
| `DB_USER` | `admin` | QuestDB Benutzername |
|
||||||
|
| `DB_PASSWORD` | `quest` | QuestDB Passwort |
|
||||||
|
| `DB_HOST` | `questdb` | QuestDB Hostname (Docker-intern) |
|
||||||
|
|
||||||
## Erweiterung
|
## Erweiterung
|
||||||
Um eine neue Börse hinzuzufügen, erstelle einfach eine neue Klasse in `src/exchanges/`, die von `BaseExchange` erbt und implementiere `fetch_latest_trades()`. Füge sie dann in `daemon.py` zur Liste hinzu.
|
|
||||||
|
Um eine neue Börse hinzuzufügen:
|
||||||
|
|
||||||
|
1. Erstelle eine neue Klasse in `src/exchanges/`, die von `BaseExchange` erbt
|
||||||
|
2. Implementiere `fetch_latest_trades()` (gibt `List[Trade]` zurück)
|
||||||
|
3. Implementiere die `name`-Property
|
||||||
|
4. Registriere in `daemon.py` unter `STREAMING_EXCHANGES` (große Daten) oder `STANDARD_EXCHANGES` (Batch)
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ from src.exchanges.boersenag import (
|
|||||||
HAMAExchange, HAMBExchange, HANAExchange, HANBExchange
|
HAMAExchange, HAMBExchange, HANAExchange, HANBExchange
|
||||||
)
|
)
|
||||||
from src.database.questdb_client import DatabaseClient
|
from src.database.questdb_client import DatabaseClient
|
||||||
|
from src.utils.validation import validate_exchange
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
@@ -67,6 +68,7 @@ STANDARD_EXCHANGES: List[Type[BaseExchange]] = [
|
|||||||
|
|
||||||
# Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert)
|
# Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert)
|
||||||
_existing_trades_cache = {}
|
_existing_trades_cache = {}
|
||||||
|
MAX_CACHE_SIZE = 50
|
||||||
|
|
||||||
def get_trade_hash(trade):
|
def get_trade_hash(trade):
|
||||||
"""Erstellt einen eindeutigen Hash für einen Trade."""
|
"""Erstellt einen eindeutigen Hash für einen Trade."""
|
||||||
@@ -75,6 +77,7 @@ def get_trade_hash(trade):
|
|||||||
|
|
||||||
def get_existing_trades_for_day(db_url, exchange_name, day):
|
def get_existing_trades_for_day(db_url, exchange_name, day):
|
||||||
"""Holt existierende Trades für einen Tag aus der DB (mit Caching)."""
|
"""Holt existierende Trades für einen Tag aus der DB (mit Caching)."""
|
||||||
|
exchange_name = validate_exchange(exchange_name)
|
||||||
cache_key = f"{exchange_name}_{day.strftime('%Y-%m-%d')}"
|
cache_key = f"{exchange_name}_{day.strftime('%Y-%m-%d')}"
|
||||||
|
|
||||||
if cache_key in _existing_trades_cache:
|
if cache_key in _existing_trades_cache:
|
||||||
@@ -109,6 +112,11 @@ def get_existing_trades_for_day(db_url, exchange_name, day):
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"Error fetching existing trades for {day}: {e}")
|
logger.warning(f"Error fetching existing trades for {day}: {e}")
|
||||||
|
|
||||||
|
# Cache-Groesse begrenzen (FIFO via dict ordering, Python 3.7+)
|
||||||
|
if len(_existing_trades_cache) >= MAX_CACHE_SIZE:
|
||||||
|
oldest_key = next(iter(_existing_trades_cache))
|
||||||
|
del _existing_trades_cache[oldest_key]
|
||||||
|
|
||||||
_existing_trades_cache[cache_key] = existing_trades
|
_existing_trades_cache[cache_key] = existing_trades
|
||||||
return existing_trades
|
return existing_trades
|
||||||
|
|
||||||
@@ -163,6 +171,7 @@ def filter_new_trades_batch(db_url, exchange_name, trades, batch_size=5000):
|
|||||||
|
|
||||||
def get_last_trade_timestamp(db_url: str, exchange_name: str) -> datetime.datetime:
|
def get_last_trade_timestamp(db_url: str, exchange_name: str) -> datetime.datetime:
|
||||||
"""Holt den Timestamp des letzten Trades für eine Exchange aus QuestDB."""
|
"""Holt den Timestamp des letzten Trades für eine Exchange aus QuestDB."""
|
||||||
|
exchange_name = validate_exchange(exchange_name)
|
||||||
query = f"trades where exchange = '{exchange_name}' latest by timestamp"
|
query = f"trades where exchange = '{exchange_name}' latest by timestamp"
|
||||||
try:
|
try:
|
||||||
response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)
|
response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)
|
||||||
|
|||||||
@@ -6,6 +6,11 @@ import requests
|
|||||||
import os
|
import os
|
||||||
import logging
|
import logging
|
||||||
from typing import Optional, Dict, Any
|
from typing import Optional, Dict, Any
|
||||||
|
from src.utils.validation import (
|
||||||
|
validate_isin, validate_exchange, validate_date,
|
||||||
|
validate_int_range, sanitize_sql_string,
|
||||||
|
validate_isin_list, validate_exchange_list,
|
||||||
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -59,7 +64,16 @@ async def get_trades(isin: str = None, days: int = 7):
|
|||||||
Gibt aggregierte Analyse aller Trades zurück (nicht einzelne Trades).
|
Gibt aggregierte Analyse aller Trades zurück (nicht einzelne Trades).
|
||||||
Nutzt vorberechnete Daten aus analytics_exchange_daily.
|
Nutzt vorberechnete Daten aus analytics_exchange_daily.
|
||||||
"""
|
"""
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 365)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-365)")
|
||||||
|
|
||||||
if isin:
|
if isin:
|
||||||
|
try:
|
||||||
|
isin = validate_isin(isin)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger ISIN-Wert")
|
||||||
# Für spezifische ISIN: hole aus trades Tabelle
|
# Für spezifische ISIN: hole aus trades Tabelle
|
||||||
query = f"""
|
query = f"""
|
||||||
select
|
select
|
||||||
@@ -109,6 +123,11 @@ async def get_summary(days: int = None):
|
|||||||
Gibt Zusammenfassung zurück. Nutzt analytics_daily_summary für total_trades.
|
Gibt Zusammenfassung zurück. Nutzt analytics_daily_summary für total_trades.
|
||||||
Optional: days Parameter für Zeitraum-basierte Zusammenfassung.
|
Optional: days Parameter für Zeitraum-basierte Zusammenfassung.
|
||||||
"""
|
"""
|
||||||
|
if days:
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 3650)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-3650)")
|
||||||
if days:
|
if days:
|
||||||
# Zeitraum-basierte Zusammenfassung
|
# Zeitraum-basierte Zusammenfassung
|
||||||
query = f"""
|
query = f"""
|
||||||
@@ -169,6 +188,11 @@ async def get_summary(days: int = None):
|
|||||||
@app.get("/api/statistics/total-trades")
|
@app.get("/api/statistics/total-trades")
|
||||||
async def get_total_trades(days: int = None):
|
async def get_total_trades(days: int = None):
|
||||||
"""Gibt Gesamtzahl aller Trades zurück (aus analytics_daily_summary). Optional: days Parameter für Zeitraum."""
|
"""Gibt Gesamtzahl aller Trades zurück (aus analytics_daily_summary). Optional: days Parameter für Zeitraum."""
|
||||||
|
if days:
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 3650)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-3650)")
|
||||||
if days:
|
if days:
|
||||||
query = f"select sum(total_trades) as total from analytics_daily_summary where timestamp >= dateadd('d', -{days}, now())"
|
query = f"select sum(total_trades) as total from analytics_daily_summary where timestamp >= dateadd('d', -{days}, now())"
|
||||||
else:
|
else:
|
||||||
@@ -201,6 +225,12 @@ async def get_custom_analytics(
|
|||||||
- exchanges: Komma-separierte Liste von Exchanges (optional)
|
- exchanges: Komma-separierte Liste von Exchanges (optional)
|
||||||
"""
|
"""
|
||||||
# Validiere Parameter
|
# Validiere Parameter
|
||||||
|
try:
|
||||||
|
date_from = validate_date(date_from)
|
||||||
|
date_to = validate_date(date_to)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges Datumsformat (erwartet: YYYY-MM-DD)")
|
||||||
|
|
||||||
valid_x_axis = ["date", "exchange", "isin"]
|
valid_x_axis = ["date", "exchange", "isin"]
|
||||||
valid_y_axis = ["volume", "trade_count", "avg_price"]
|
valid_y_axis = ["volume", "trade_count", "avg_price"]
|
||||||
valid_group_by = ["exchange", "sector", "date"]
|
valid_group_by = ["exchange", "sector", "date"]
|
||||||
@@ -212,6 +242,13 @@ async def get_custom_analytics(
|
|||||||
if group_by not in valid_group_by:
|
if group_by not in valid_group_by:
|
||||||
raise HTTPException(status_code=400, detail=f"Invalid group_by. Must be one of: {valid_group_by}")
|
raise HTTPException(status_code=400, detail=f"Invalid group_by. Must be one of: {valid_group_by}")
|
||||||
|
|
||||||
|
validated_exchanges = None
|
||||||
|
if exchanges:
|
||||||
|
try:
|
||||||
|
validated_exchanges = validate_exchange_list(exchanges)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger Exchange-Name in der Liste")
|
||||||
|
|
||||||
# Für Sektor-Gruppierung: direkter JOIN mit metadata (nicht vorberechnet)
|
# Für Sektor-Gruppierung: direkter JOIN mit metadata (nicht vorberechnet)
|
||||||
if group_by == "sector":
|
if group_by == "sector":
|
||||||
y_axis_map = {
|
y_axis_map = {
|
||||||
@@ -233,8 +270,8 @@ async def get_custom_analytics(
|
|||||||
and t.timestamp <= '{date_to}'
|
and t.timestamp <= '{date_to}'
|
||||||
"""
|
"""
|
||||||
|
|
||||||
if exchanges:
|
if validated_exchanges:
|
||||||
exchange_list = ",".join([f"'{e.strip()}'" for e in exchanges.split(",")])
|
exchange_list = ",".join([f"'{e}'" for e in validated_exchanges])
|
||||||
query += f" and t.exchange in ({exchange_list})"
|
query += f" and t.exchange in ({exchange_list})"
|
||||||
|
|
||||||
query += f" group by date_trunc('day', t.timestamp), coalesce(m.sector, 'Unbekannt') order by x_value asc, group_value asc"
|
query += f" group by date_trunc('day', t.timestamp), coalesce(m.sector, 'Unbekannt') order by x_value asc, group_value asc"
|
||||||
@@ -252,12 +289,9 @@ async def get_custom_analytics(
|
|||||||
|
|
||||||
# Nutze vorberechnete Daten aus analytics_custom
|
# Nutze vorberechnete Daten aus analytics_custom
|
||||||
exchange_filter = "all"
|
exchange_filter = "all"
|
||||||
if exchanges:
|
if validated_exchanges:
|
||||||
# Wenn mehrere Exchanges angegeben, müssen wir kombinieren
|
if len(validated_exchanges) == 1:
|
||||||
# Für jetzt: nutze nur wenn ein Exchange angegeben ist
|
exchange_filter = validated_exchanges[0]
|
||||||
exchange_list = [e.strip() for e in exchanges.split(",")]
|
|
||||||
if len(exchange_list) == 1:
|
|
||||||
exchange_filter = exchange_list[0]
|
|
||||||
else:
|
else:
|
||||||
# Bei mehreren Exchanges: gib Fehler zurück, da dies nicht vorberechnet wird
|
# Bei mehreren Exchanges: gib Fehler zurück, da dies nicht vorberechnet wird
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
@@ -317,6 +351,10 @@ async def get_moving_average(days: int = 7, exchange: str = None):
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
if exchange:
|
if exchange:
|
||||||
|
try:
|
||||||
|
exchange = validate_exchange(exchange)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger Exchange-Name")
|
||||||
query += f" and exchange = '{exchange}'"
|
query += f" and exchange = '{exchange}'"
|
||||||
|
|
||||||
query += " order by date asc, exchange asc"
|
query += " order by date asc, exchange asc"
|
||||||
@@ -393,6 +431,10 @@ async def get_stock_trends(days: int = 7, limit: int = 20):
|
|||||||
"""
|
"""
|
||||||
if days not in [7, 30, 42, 69, 180, 365]:
|
if days not in [7, 30, 42, 69, 180, 365]:
|
||||||
raise HTTPException(status_code=400, detail="Invalid days parameter. Must be one of: 7, 30, 42, 69, 180, 365")
|
raise HTTPException(status_code=400, detail="Invalid days parameter. Must be one of: 7, 30, 42, 69, 180, 365")
|
||||||
|
try:
|
||||||
|
limit = validate_int_range(limit, 1, 1000)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger limit-Parameter (1-1000)")
|
||||||
|
|
||||||
query = f"""
|
query = f"""
|
||||||
select
|
select
|
||||||
@@ -423,7 +465,40 @@ async def get_analytics(
|
|||||||
continents: str = None
|
continents: str = None
|
||||||
):
|
):
|
||||||
"""Analytics Endpunkt für Report Builder"""
|
"""Analytics Endpunkt für Report Builder"""
|
||||||
|
# Validiere optionale Parameter
|
||||||
|
if date_from:
|
||||||
|
try:
|
||||||
|
date_from = validate_date(date_from)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges date_from Format (erwartet: YYYY-MM-DD)")
|
||||||
|
if date_to:
|
||||||
|
try:
|
||||||
|
date_to = validate_date(date_to)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges date_to Format (erwartet: YYYY-MM-DD)")
|
||||||
|
|
||||||
|
validated_isins = None
|
||||||
|
if isins:
|
||||||
|
try:
|
||||||
|
validated_isins = validate_isin_list(isins)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger ISIN-Wert in der Liste")
|
||||||
|
|
||||||
|
sanitized_continents = None
|
||||||
|
if continents:
|
||||||
|
sanitized_continents = [sanitize_sql_string(c.strip(), max_length=50) for c in continents.split(",") if c.strip()]
|
||||||
|
|
||||||
composite_keys = ["exchange_continent", "exchange_sector"]
|
composite_keys = ["exchange_continent", "exchange_sector"]
|
||||||
|
valid_metrics = ["volume", "count", "avg_price", "all"]
|
||||||
|
valid_groups = ["day", "month", "exchange", "isin", "name", "continent", "sector", "exchange_continent", "exchange_sector"]
|
||||||
|
|
||||||
|
if metric not in valid_metrics:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger metric-Parameter. Erlaubt: {valid_metrics}")
|
||||||
|
if group_by not in valid_groups:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger group_by-Parameter. Erlaubt: {valid_groups}")
|
||||||
|
if sub_group_by and sub_group_by not in valid_groups:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger sub_group_by-Parameter. Erlaubt: {valid_groups}")
|
||||||
|
|
||||||
needs_metadata = any([
|
needs_metadata = any([
|
||||||
group_by in ["name", "continent", "sector"] + composite_keys,
|
group_by in ["name", "continent", "sector"] + composite_keys,
|
||||||
sub_group_by in ["name", "continent", "sector"] + composite_keys,
|
sub_group_by in ["name", "continent", "sector"] + composite_keys,
|
||||||
@@ -473,12 +548,12 @@ async def get_analytics(
|
|||||||
if date_to:
|
if date_to:
|
||||||
query += f" and {t_prefix}timestamp <= '{date_to}'"
|
query += f" and {t_prefix}timestamp <= '{date_to}'"
|
||||||
|
|
||||||
if isins:
|
if validated_isins:
|
||||||
isins_list = ",".join([f"'{i.strip()}'" for i in isins.split(",")])
|
isins_list = ",".join([f"'{i}'" for i in validated_isins])
|
||||||
query += f" and {t_prefix}isin in ({isins_list})"
|
query += f" and {t_prefix}isin in ({isins_list})"
|
||||||
|
|
||||||
if continents and needs_metadata:
|
if sanitized_continents and needs_metadata:
|
||||||
cont_list = ",".join([f"'{c.strip()}'" for c in continents.split(",")])
|
cont_list = ",".join([f"'{c}'" for c in sanitized_continents])
|
||||||
query += f" and {m_prefix}continent in ({cont_list})"
|
query += f" and {m_prefix}continent in ({cont_list})"
|
||||||
|
|
||||||
query += f" group by {selected_group}"
|
query += f" group by {selected_group}"
|
||||||
@@ -493,7 +568,8 @@ async def get_analytics(
|
|||||||
@app.get("/api/metadata/search")
|
@app.get("/api/metadata/search")
|
||||||
async def search_metadata(q: str):
|
async def search_metadata(q: str):
|
||||||
"""Case-insensitive search for ISIN or Name"""
|
"""Case-insensitive search for ISIN or Name"""
|
||||||
query = f"select isin, name from metadata where isin ilike '%{q}%' or name ilike '%{q}%' limit 10"
|
q_safe = sanitize_sql_string(q, max_length=100)
|
||||||
|
query = f"select isin, name from metadata where isin ilike '%{q_safe}%' or name ilike '%{q_safe}%' limit 10"
|
||||||
data = query_questdb(query)
|
data = query_questdb(query)
|
||||||
return format_questdb_response(data)
|
return format_questdb_response(data)
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import os
|
|||||||
import requests
|
import requests
|
||||||
from typing import Dict, List, Tuple, Optional
|
from typing import Dict, List, Tuple, Optional
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
|
from src.utils.validation import validate_table_name, validate_exchange
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
@@ -499,6 +500,10 @@ class AnalyticsWorker:
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
if exchange_filter:
|
if exchange_filter:
|
||||||
|
try:
|
||||||
|
exchange_filter = validate_exchange(exchange_filter)
|
||||||
|
except ValueError:
|
||||||
|
continue
|
||||||
query += f" and exchange = '{exchange_filter}'"
|
query += f" and exchange = '{exchange_filter}'"
|
||||||
|
|
||||||
query += f" group by date_trunc('day', timestamp), {group_by_field}"
|
query += f" group by date_trunc('day', timestamp), {group_by_field}"
|
||||||
@@ -764,6 +769,7 @@ class AnalyticsWorker:
|
|||||||
|
|
||||||
def get_existing_dates(self, table_name: str) -> set:
|
def get_existing_dates(self, table_name: str) -> set:
|
||||||
"""Holt alle bereits berechneten Daten aus einer Analytics-Tabelle"""
|
"""Holt alle bereits berechneten Daten aus einer Analytics-Tabelle"""
|
||||||
|
table_name = validate_table_name(table_name)
|
||||||
query = f"select distinct date_trunc('day', timestamp) as date from {table_name}"
|
query = f"select distinct date_trunc('day', timestamp) as date from {table_name}"
|
||||||
data = self.query_questdb(query)
|
data = self.query_questdb(query)
|
||||||
if not data:
|
if not data:
|
||||||
@@ -875,12 +881,13 @@ class AnalyticsWorker:
|
|||||||
|
|
||||||
for table in tables:
|
for table in tables:
|
||||||
try:
|
try:
|
||||||
|
table = validate_table_name(table)
|
||||||
# QuestDB DELETE syntax
|
# QuestDB DELETE syntax
|
||||||
delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'"
|
delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'"
|
||||||
response = requests.get(
|
response = requests.get(
|
||||||
f"{self.questdb_url}/exec",
|
f"{self.db_url}/exec",
|
||||||
params={'query': delete_query},
|
params={'query': delete_query},
|
||||||
auth=self.auth,
|
auth=DB_AUTH,
|
||||||
timeout=30
|
timeout=30
|
||||||
)
|
)
|
||||||
if response.status_code == 200:
|
if response.status_code == 200:
|
||||||
|
|||||||
@@ -1,8 +1,11 @@
|
|||||||
import requests
|
import requests
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from typing import List
|
from typing import List
|
||||||
from ..exchanges.base import Trade
|
from ..exchanges.base import Trade
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
class DatabaseClient:
|
class DatabaseClient:
|
||||||
def __init__(self, host: str = "localhost", port: int = 9000, user: str = None, password: str = None):
|
def __init__(self, host: str = "localhost", port: int = 9000, user: str = None, password: str = None):
|
||||||
self.host = host
|
self.host = host
|
||||||
@@ -15,7 +18,7 @@ class DatabaseClient:
|
|||||||
return
|
return
|
||||||
|
|
||||||
total_trades = len(trades)
|
total_trades = len(trades)
|
||||||
print(f"Saving {total_trades} trades to QuestDB in batches of {batch_size}...")
|
logger.info(f"Speichere {total_trades} Trades in QuestDB (Batches von {batch_size})...")
|
||||||
|
|
||||||
for i in range(0, total_trades, batch_size):
|
for i in range(0, total_trades, batch_size):
|
||||||
batch = trades[i:i + batch_size]
|
batch = trades[i:i + batch_size]
|
||||||
@@ -31,7 +34,7 @@ class DatabaseClient:
|
|||||||
f"{int(trade.timestamp.timestamp() * 1e9)}"
|
f"{int(trade.timestamp.timestamp() * 1e9)}"
|
||||||
lines.append(line)
|
lines.append(line)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Error formating trade {trade}: {e}")
|
logger.error(f"Fehler beim Formatieren von Trade {trade}: {e}")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if not lines:
|
if not lines:
|
||||||
@@ -47,12 +50,12 @@ class DatabaseClient:
|
|||||||
auth=self.auth
|
auth=self.auth
|
||||||
)
|
)
|
||||||
if response.status_code not in [204, 200]:
|
if response.status_code not in [204, 200]:
|
||||||
print(f"Error saving batch {i//batch_size + 1} to QuestDB: {response.text}")
|
logger.error(f"Fehler beim Speichern von Batch {i//batch_size + 1}: {response.text}")
|
||||||
else:
|
else:
|
||||||
print(f"Saved batch {i//batch_size + 1} ({len(batch)} trades)")
|
logger.info(f"Batch {i//batch_size + 1} gespeichert ({len(batch)} Trades)")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Could not connect to QuestDB at {self.url}: {e}")
|
logger.error(f"Verbindung zu QuestDB fehlgeschlagen ({self.url}): {e}")
|
||||||
# Fallback: print to console or save to file
|
# Fallback: in Datei speichern
|
||||||
self._fallback_save(batch)
|
self._fallback_save(batch)
|
||||||
|
|
||||||
def _fallback_save(self, trades: List[Trade]):
|
def _fallback_save(self, trades: List[Trade]):
|
||||||
|
|||||||
@@ -8,11 +8,14 @@ URL-Format: https://cld42.boersenag.de/m13data/data/Mifir13DelayedData_{MIC}_{SE
|
|||||||
|
|
||||||
import requests
|
import requests
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
import re
|
import re
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Rate-Limiting Konfiguration
|
# Rate-Limiting Konfiguration
|
||||||
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
||||||
|
|
||||||
@@ -176,9 +179,9 @@ class BoersenagBase(BaseExchange):
|
|||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[{self.name}] HTTP error: {e}")
|
logger.error(f"[{self.name}] HTTP error: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -265,9 +268,9 @@ class BoersenagBase(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Generiere mögliche URLs
|
# Generiere mögliche URLs
|
||||||
urls = self._generate_file_urls(target_date)
|
urls = self._generate_file_urls(target_date)
|
||||||
@@ -281,7 +284,7 @@ class BoersenagBase(BaseExchange):
|
|||||||
if trades:
|
if trades:
|
||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful += 1
|
successful += 1
|
||||||
print(f"[{self.name}] Found {len(trades)} trades from: {url.split('/')[-1]}")
|
logger.info(f"[{self.name}] Found {len(trades)} trades from: {url.split('/')[-1]}")
|
||||||
# Bei Erfolg müssen wir nicht alle anderen URLs probieren
|
# Bei Erfolg müssen wir nicht alle anderen URLs probieren
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -293,7 +296,7 @@ class BoersenagBase(BaseExchange):
|
|||||||
if i > 20 and successful == 0:
|
if i > 20 and successful == 0:
|
||||||
break
|
break
|
||||||
|
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -3,11 +3,14 @@ import gzip
|
|||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
from bs4 import BeautifulSoup
|
from bs4 import BeautifulSoup
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Rate-Limiting
|
# Rate-Limiting
|
||||||
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
||||||
|
|
||||||
@@ -79,10 +82,10 @@ class GettexExchange(BaseExchange):
|
|||||||
url = f"https://www.gettex.de/fileadmin/posttrade-data/{filename}"
|
url = f"https://www.gettex.de/fileadmin/posttrade-data/{filename}"
|
||||||
files.append({'filename': filename, 'url': url})
|
files.append({'filename': filename, 'url': url})
|
||||||
|
|
||||||
print(f"[GETTEX] Found {len(files)} files on page")
|
logger.info(f"[GETTEX] Found {len(files)} files on page")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[GETTEX] Error fetching page: {e}")
|
logger.error(f"[GETTEX] Error fetching page: {e}")
|
||||||
|
|
||||||
return files
|
return files
|
||||||
|
|
||||||
@@ -148,11 +151,11 @@ class GettexExchange(BaseExchange):
|
|||||||
date_str = parts[1] # YYYYMMDD
|
date_str = parts[1] # YYYYMMDD
|
||||||
|
|
||||||
if not date_str:
|
if not date_str:
|
||||||
print(f"[GETTEX] WARNING: Could not extract date from filename: {filename}")
|
logger.warning(f"[GETTEX] Could not extract date from filename: {filename}")
|
||||||
|
|
||||||
# Debug: Zeige erste Zeile
|
# Debug: Zeige erste Zeile
|
||||||
if lines and len(lines) > 0:
|
if lines and len(lines) > 0:
|
||||||
print(f"[GETTEX] First line sample: {lines[0][:100]}")
|
logger.debug(f"[GETTEX] First line sample: {lines[0][:100]}")
|
||||||
|
|
||||||
# Gettex CSV hat KEINEN Header!
|
# Gettex CSV hat KEINEN Header!
|
||||||
# Format: ISIN,Zeit,Währung,Preis,Menge
|
# Format: ISIN,Zeit,Währung,Preis,Menge
|
||||||
@@ -166,24 +169,24 @@ class GettexExchange(BaseExchange):
|
|||||||
if trade:
|
if trade:
|
||||||
trades.append(trade)
|
trades.append(trade)
|
||||||
else:
|
else:
|
||||||
if i < 3: # Zeige nur erste paar Fehler
|
if i < 3:
|
||||||
print(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
logger.debug(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
parse_errors += 1
|
parse_errors += 1
|
||||||
if i < 3:
|
if i < 3:
|
||||||
print(f"[GETTEX] Exception parsing line {i+1}: {e}, line: {line[:80]}")
|
logger.debug(f"[GETTEX] Exception parsing line {i+1}: {e}, line: {line[:80]}")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if trades:
|
if trades:
|
||||||
print(f"[GETTEX] Parsed {len(trades)} trades from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
logger.info(f"[GETTEX] Parsed {len(trades)} trades from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
||||||
elif len(lines) > 0:
|
elif len(lines) > 0:
|
||||||
print(f"[GETTEX] No trades parsed from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
logger.warning(f"[GETTEX] No trades parsed from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[GETTEX] HTTP error downloading {filename}: {e}")
|
logger.error(f"[GETTEX] HTTP error downloading {filename}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[GETTEX] Error downloading {filename}: {e}")
|
logger.error(f"[GETTEX] Error downloading {filename}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -406,9 +409,9 @@ class GettexExchange(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Versuche zuerst, Dateien von der Webseite zu laden
|
# Versuche zuerst, Dateien von der Webseite zu laden
|
||||||
page_files = self._get_file_list_from_page()
|
page_files = self._get_file_list_from_page()
|
||||||
@@ -434,10 +437,10 @@ class GettexExchange(BaseExchange):
|
|||||||
hour = int(parts[2])
|
hour = int(parts[2])
|
||||||
if hour < 3:
|
if hour < 3:
|
||||||
target_files.append(f)
|
target_files.append(f)
|
||||||
except:
|
except (ValueError, IndexError):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
print(f"[{self.name}] Found {len(target_files)} files for target date from page")
|
logger.info(f"[{self.name}] Found {len(target_files)} files for target date from page")
|
||||||
|
|
||||||
# Lade Dateien von der Webseite (mit Rate-Limiting)
|
# Lade Dateien von der Webseite (mit Rate-Limiting)
|
||||||
for i, f in enumerate(target_files):
|
for i, f in enumerate(target_files):
|
||||||
@@ -450,9 +453,9 @@ class GettexExchange(BaseExchange):
|
|||||||
|
|
||||||
# Fallback: Versuche erwartete Dateinamen
|
# Fallback: Versuche erwartete Dateinamen
|
||||||
if not all_trades:
|
if not all_trades:
|
||||||
print(f"[{self.name}] No files from page, trying generated filenames...")
|
logger.info(f"[{self.name}] No files from page, trying generated filenames...")
|
||||||
expected_files = self._generate_expected_files(target_date)
|
expected_files = self._generate_expected_files(target_date)
|
||||||
print(f"[{self.name}] Trying {len(expected_files)} potential files")
|
logger.info(f"[{self.name}] Trying {len(expected_files)} potential files")
|
||||||
|
|
||||||
successful_files = 0
|
successful_files = 0
|
||||||
for filename in expected_files:
|
for filename in expected_files:
|
||||||
@@ -461,9 +464,9 @@ class GettexExchange(BaseExchange):
|
|||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful_files += 1
|
successful_files += 1
|
||||||
|
|
||||||
print(f"[{self.name}] Successfully downloaded {successful_files} files")
|
logger.info(f"[{self.name}] Successfully downloaded {successful_files} files")
|
||||||
|
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
|
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|
||||||
@@ -506,12 +509,12 @@ class GettexExchange(BaseExchange):
|
|||||||
continue
|
continue
|
||||||
|
|
||||||
if trades:
|
if trades:
|
||||||
print(f"[{self.name}] Parsed {len(trades)} trades from {filename}")
|
logger.info(f"[{self.name}] Parsed {len(trades)} trades from {filename}")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[{self.name}] HTTP error downloading {url}: {e}")
|
logger.error(f"[{self.name}] HTTP error downloading {url}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|||||||
@@ -1,10 +1,13 @@
|
|||||||
import requests
|
import requests
|
||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
|
import logging
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from typing import List
|
from typing import List
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
class LSExchange(BaseExchange):
|
class LSExchange(BaseExchange):
|
||||||
@property
|
@property
|
||||||
def name(self) -> str:
|
def name(self) -> str:
|
||||||
@@ -52,8 +55,9 @@ class LSExchange(BaseExchange):
|
|||||||
quantity=quantity,
|
quantity=quantity,
|
||||||
timestamp=timestamp
|
timestamp=timestamp
|
||||||
))
|
))
|
||||||
except Exception:
|
except (ValueError, KeyError) as e:
|
||||||
|
logger.debug(f"Fehler beim Parsen einer LS-Zeile: {e}")
|
||||||
continue
|
continue
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Error fetching LS data from {url}: {e}")
|
logger.error(f"Fehler beim Abrufen von LS-Daten von {url}: {e}")
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|||||||
@@ -3,11 +3,14 @@ import gzip
|
|||||||
import json
|
import json
|
||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
from bs4 import BeautifulSoup
|
from bs4 import BeautifulSoup
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Browser User-Agent (Vollständiger Browser-Fingerprint für Stuttgart)
|
# Browser User-Agent (Vollständiger Browser-Fingerprint für Stuttgart)
|
||||||
HEADERS = {
|
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',
|
'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',
|
||||||
@@ -84,7 +87,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
files = self._generate_expected_urls()
|
files = self._generate_expected_urls()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error fetching page: {e}")
|
logger.error(f"[STU] Error fetching page: {e}")
|
||||||
files = self._generate_expected_urls()
|
files = self._generate_expected_urls()
|
||||||
|
|
||||||
return files
|
return files
|
||||||
@@ -194,13 +197,13 @@ class StuttgartExchange(BaseExchange):
|
|||||||
if trade:
|
if trade:
|
||||||
trades.append(trade)
|
trades.append(trade)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Could not parse {url}: {e}")
|
logger.error(f"[STU] Could not parse {url}: {e}")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[STU] HTTP error downloading {url}: {e}")
|
logger.error(f"[STU] HTTP error downloading {url}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error downloading {url}: {e}")
|
logger.error(f"[STU] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -275,7 +278,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error parsing JSON record: {e}")
|
logger.debug(f"[STU] Error parsing JSON record: {e}")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def _parse_csv_row(self, row: dict) -> Optional[Trade]:
|
def _parse_csv_row(self, row: dict) -> Optional[Trade]:
|
||||||
@@ -331,7 +334,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error parsing CSV row: {e}")
|
logger.debug(f"[STU] Error parsing CSV row: {e}")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def _get_last_trading_day(self, from_date) -> datetime.date:
|
def _get_last_trading_day(self, from_date) -> datetime.date:
|
||||||
@@ -365,13 +368,13 @@ class StuttgartExchange(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Download-Links holen
|
# Download-Links holen
|
||||||
all_links = self._get_download_links()
|
all_links = self._get_download_links()
|
||||||
print(f"[{self.name}] Found {len(all_links)} potential download links")
|
logger.info(f"[{self.name}] Found {len(all_links)} potential download links")
|
||||||
|
|
||||||
# Nach Datum filtern
|
# Nach Datum filtern
|
||||||
target_links = self._filter_files_for_date(all_links, target_date)
|
target_links = self._filter_files_for_date(all_links, target_date)
|
||||||
@@ -380,7 +383,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
# Fallback: Versuche alle Links
|
# Fallback: Versuche alle Links
|
||||||
target_links = all_links
|
target_links = all_links
|
||||||
|
|
||||||
print(f"[{self.name}] Trying {len(target_links)} files for target date")
|
logger.info(f"[{self.name}] Trying {len(target_links)} files for target date")
|
||||||
|
|
||||||
# Dateien herunterladen und parsen
|
# Dateien herunterladen und parsen
|
||||||
successful = 0
|
successful = 0
|
||||||
@@ -389,9 +392,9 @@ class StuttgartExchange(BaseExchange):
|
|||||||
if trades:
|
if trades:
|
||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful += 1
|
successful += 1
|
||||||
print(f"[{self.name}] Parsed {len(trades)} trades from {url}")
|
logger.info(f"[{self.name}] Parsed {len(trades)} trades from {url}")
|
||||||
|
|
||||||
print(f"[{self.name}] Successfully processed {successful} files")
|
logger.info(f"[{self.name}] Successfully processed {successful} files")
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
|
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|||||||
0
src/utils/__init__.py
Normal file
0
src/utils/__init__.py
Normal file
124
src/utils/validation.py
Normal file
124
src/utils/validation.py
Normal file
@@ -0,0 +1,124 @@
|
|||||||
|
"""
|
||||||
|
Zentrale Validierungs- und Sanitisierungsfunktionen fuer SQL-Queries.
|
||||||
|
|
||||||
|
Da QuestDB's HTTP API keine parametrisierten Queries unterstuetzt,
|
||||||
|
muessen alle Werte vor der Interpolation validiert/bereinigt werden.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import re
|
||||||
|
import datetime
|
||||||
|
|
||||||
|
|
||||||
|
# Gueltige Exchange-Namen (alle registrierten Boersen)
|
||||||
|
VALID_EXCHANGES = {
|
||||||
|
"EIX", "LS", "XETRA", "FRA", "QUOTRIX", "GETTEX", "STU",
|
||||||
|
"DUSA", "DUSB", "DUSC", "DUSD",
|
||||||
|
"HAMA", "HAMB", "HANA", "HANB",
|
||||||
|
"NONE", # Platzhalter fuer leere Analytics-Eintraege
|
||||||
|
}
|
||||||
|
|
||||||
|
# Gueltige Analytics-Tabellennamen
|
||||||
|
VALID_TABLES = {
|
||||||
|
"analytics_custom", "analytics_exchange_daily",
|
||||||
|
"analytics_daily_summary", "analytics_volume_changes",
|
||||||
|
"analytics_stock_trends", "trades", "metadata",
|
||||||
|
}
|
||||||
|
|
||||||
|
# ISIN-Format: 2 Buchstaben Laendercode + 9 alphanumerische Zeichen + 1 Pruefziffer
|
||||||
|
_ISIN_PATTERN = re.compile(r'^[A-Z]{2}[A-Z0-9]{9}[0-9]$')
|
||||||
|
|
||||||
|
# Datumsformat: YYYY-MM-DD
|
||||||
|
_DATE_PATTERN = re.compile(r'^\d{4}-\d{2}-\d{2}$')
|
||||||
|
|
||||||
|
# Gefaehrliche SQL-Fragmente
|
||||||
|
_SQL_DANGEROUS = re.compile(r'(--|/\*|\*/|;)')
|
||||||
|
|
||||||
|
|
||||||
|
def validate_isin(value: str) -> str:
|
||||||
|
"""Validiert einen ISIN-Wert. Wirft ValueError bei ungueltigem Format."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"ISIN muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().upper()
|
||||||
|
if not _ISIN_PATTERN.match(value):
|
||||||
|
raise ValueError(f"Ungueltiges ISIN-Format: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_exchange(value: str) -> str:
|
||||||
|
"""Validiert einen Exchange-Namen gegen die Whitelist."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Exchange muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().upper()
|
||||||
|
if value not in VALID_EXCHANGES:
|
||||||
|
raise ValueError(f"Unbekannte Exchange: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_date(value: str) -> str:
|
||||||
|
"""Validiert ein Datum im Format YYYY-MM-DD."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Datum muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip()
|
||||||
|
if not _DATE_PATTERN.match(value):
|
||||||
|
raise ValueError(f"Ungueltiges Datumsformat: {value!r} (erwartet: YYYY-MM-DD)")
|
||||||
|
# Pruefe ob es ein gueltiges Datum ist
|
||||||
|
try:
|
||||||
|
datetime.date.fromisoformat(value)
|
||||||
|
except ValueError:
|
||||||
|
raise ValueError(f"Ungueltiges Datum: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_int_range(value: int, min_val: int = 0, max_val: int = 10000) -> int:
|
||||||
|
"""Validiert einen Integer-Wert innerhalb eines Bereichs."""
|
||||||
|
value = int(value)
|
||||||
|
if value < min_val or value > max_val:
|
||||||
|
raise ValueError(f"Wert {value} ausserhalb des erlaubten Bereichs [{min_val}, {max_val}]")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def sanitize_sql_string(value: str, max_length: int = 200) -> str:
|
||||||
|
"""
|
||||||
|
Bereinigt einen String fuer die Verwendung in SQL-Queries.
|
||||||
|
Escapet Single-Quotes und entfernt gefaehrliche SQL-Fragmente.
|
||||||
|
"""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Wert muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value[:max_length]
|
||||||
|
# Single-Quotes escapen (SQL-Standard)
|
||||||
|
value = value.replace("'", "''")
|
||||||
|
# Gefaehrliche SQL-Fragmente entfernen
|
||||||
|
value = _SQL_DANGEROUS.sub('', value)
|
||||||
|
# Backslashes entfernen
|
||||||
|
value = value.replace('\\', '')
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_isin_list(value: str) -> list:
|
||||||
|
"""Validiert eine komma-separierte Liste von ISINs."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError("ISIN-Liste muss ein String sein")
|
||||||
|
items = [item.strip() for item in value.split(",") if item.strip()]
|
||||||
|
if not items:
|
||||||
|
raise ValueError("ISIN-Liste ist leer")
|
||||||
|
return [validate_isin(item) for item in items]
|
||||||
|
|
||||||
|
|
||||||
|
def validate_exchange_list(value: str) -> list:
|
||||||
|
"""Validiert eine komma-separierte Liste von Exchange-Namen."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError("Exchange-Liste muss ein String sein")
|
||||||
|
items = [item.strip() for item in value.split(",") if item.strip()]
|
||||||
|
if not items:
|
||||||
|
raise ValueError("Exchange-Liste ist leer")
|
||||||
|
return [validate_exchange(item) for item in items]
|
||||||
|
|
||||||
|
|
||||||
|
def validate_table_name(value: str) -> str:
|
||||||
|
"""Validiert einen Tabellennamen gegen die Whitelist."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Tabellenname muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().lower()
|
||||||
|
if value not in VALID_TABLES:
|
||||||
|
raise ValueError(f"Unbekannter Tabellenname: {value!r}")
|
||||||
|
return value
|
||||||
Reference in New Issue
Block a user