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 . .
|
||||
|
||||
ENV PYTHONPATH=/app
|
||||
|
||||
CMD ["python", "dashboard/server.py"]
|
||||
|
||||
171
README.md
171
README.md
@@ -1,20 +1,140 @@
|
||||
# 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
|
||||
- **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
|
||||
- `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).
|
||||
- `daemon.py`: Der Orchestrator, der die Daten abruft und speichert.
|
||||
|
||||
Fünf Microservices, orchestriert via Docker Compose:
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────────────────────────────┐
|
||||
│ 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
|
||||
|
||||
### 1. QuestDB (Timeseries DB) starten
|
||||
Am einfachsten via Docker Compose:
|
||||
### 1. QuestDB starten (via Docker Compose)
|
||||
```bash
|
||||
docker-compose up -d
|
||||
```
|
||||
@@ -25,26 +145,35 @@ QuestDB ist dann unter `http://localhost:9000` erreichbar.
|
||||
pip install -r requirements.txt
|
||||
```
|
||||
|
||||
### 3. Systemd Service einrichten
|
||||
Kopiere die Dateien nach `/etc/systemd/system/`:
|
||||
### 3. Manuell starten
|
||||
```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
|
||||
sudo cp systemd/trading-daemon.service /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 enable --now trading-daemon.timer
|
||||
```
|
||||
|
||||
### 4. Manuell testen
|
||||
```bash
|
||||
python3 daemon.py
|
||||
```
|
||||
## Umgebungsvariablen
|
||||
|
||||
| Variable | Default | Beschreibung |
|
||||
|----------|---------|-------------|
|
||||
| `DB_USER` | `admin` | QuestDB Benutzername |
|
||||
| `DB_PASSWORD` | `quest` | QuestDB Passwort |
|
||||
| `DB_HOST` | `questdb` | QuestDB Hostname (Docker-intern) |
|
||||
|
||||
## 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)
|
||||
|
||||
11
daemon.py
11
daemon.py
@@ -17,6 +17,7 @@ from src.exchanges.boersenag import (
|
||||
HAMAExchange, HAMBExchange, HANAExchange, HANBExchange
|
||||
)
|
||||
from src.database.questdb_client import DatabaseClient
|
||||
from src.utils.validation import validate_exchange
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
@@ -67,6 +68,7 @@ STANDARD_EXCHANGES: List[Type[BaseExchange]] = [
|
||||
|
||||
# Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert)
|
||||
_existing_trades_cache = {}
|
||||
MAX_CACHE_SIZE = 50
|
||||
|
||||
def get_trade_hash(trade):
|
||||
"""Erstellt einen eindeutigen Hash für einen Trade."""
|
||||
@@ -75,8 +77,9 @@ def get_trade_hash(trade):
|
||||
|
||||
def get_existing_trades_for_day(db_url, exchange_name, day):
|
||||
"""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')}"
|
||||
|
||||
|
||||
if cache_key in _existing_trades_cache:
|
||||
return _existing_trades_cache[cache_key]
|
||||
|
||||
@@ -109,6 +112,11 @@ def get_existing_trades_for_day(db_url, exchange_name, day):
|
||||
except Exception as 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
|
||||
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:
|
||||
"""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"
|
||||
try:
|
||||
response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)
|
||||
|
||||
@@ -6,6 +6,11 @@ import requests
|
||||
import os
|
||||
import logging
|
||||
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__)
|
||||
|
||||
@@ -59,10 +64,19 @@ async def get_trades(isin: str = None, days: int = 7):
|
||||
Gibt aggregierte Analyse aller Trades zurück (nicht einzelne Trades).
|
||||
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:
|
||||
try:
|
||||
isin = validate_isin(isin)
|
||||
except ValueError:
|
||||
raise HTTPException(status_code=400, detail="Ungueltiger ISIN-Wert")
|
||||
# Für spezifische ISIN: hole aus trades Tabelle
|
||||
query = f"""
|
||||
select
|
||||
select
|
||||
date_trunc('day', timestamp) as date,
|
||||
count(*) as trade_count,
|
||||
sum(price * quantity) as volume,
|
||||
@@ -109,6 +123,11 @@ async def get_summary(days: int = None):
|
||||
Gibt Zusammenfassung zurück. Nutzt analytics_daily_summary für total_trades.
|
||||
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:
|
||||
# Zeitraum-basierte Zusammenfassung
|
||||
query = f"""
|
||||
@@ -169,6 +188,11 @@ async def get_summary(days: int = None):
|
||||
@app.get("/api/statistics/total-trades")
|
||||
async def get_total_trades(days: int = None):
|
||||
"""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:
|
||||
query = f"select sum(total_trades) as total from analytics_daily_summary where timestamp >= dateadd('d', -{days}, now())"
|
||||
else:
|
||||
@@ -201,17 +225,30 @@ async def get_custom_analytics(
|
||||
- exchanges: Komma-separierte Liste von Exchanges (optional)
|
||||
"""
|
||||
# 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_y_axis = ["volume", "trade_count", "avg_price"]
|
||||
valid_group_by = ["exchange", "sector", "date"]
|
||||
|
||||
|
||||
if x_axis not in valid_x_axis:
|
||||
raise HTTPException(status_code=400, detail=f"Invalid x_axis. Must be one of: {valid_x_axis}")
|
||||
if y_axis not in valid_y_axis:
|
||||
raise HTTPException(status_code=400, detail=f"Invalid y_axis. Must be one of: {valid_y_axis}")
|
||||
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}")
|
||||
|
||||
|
||||
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)
|
||||
if group_by == "sector":
|
||||
y_axis_map = {
|
||||
@@ -233,8 +270,8 @@ async def get_custom_analytics(
|
||||
and t.timestamp <= '{date_to}'
|
||||
"""
|
||||
|
||||
if exchanges:
|
||||
exchange_list = ",".join([f"'{e.strip()}'" for e in exchanges.split(",")])
|
||||
if validated_exchanges:
|
||||
exchange_list = ",".join([f"'{e}'" for e in validated_exchanges])
|
||||
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"
|
||||
@@ -252,12 +289,9 @@ async def get_custom_analytics(
|
||||
|
||||
# Nutze vorberechnete Daten aus analytics_custom
|
||||
exchange_filter = "all"
|
||||
if exchanges:
|
||||
# Wenn mehrere Exchanges angegeben, müssen wir kombinieren
|
||||
# Für jetzt: nutze nur wenn ein Exchange angegeben ist
|
||||
exchange_list = [e.strip() for e in exchanges.split(",")]
|
||||
if len(exchange_list) == 1:
|
||||
exchange_filter = exchange_list[0]
|
||||
if validated_exchanges:
|
||||
if len(validated_exchanges) == 1:
|
||||
exchange_filter = validated_exchanges[0]
|
||||
else:
|
||||
# Bei mehreren Exchanges: gib Fehler zurück, da dies nicht vorberechnet wird
|
||||
raise HTTPException(
|
||||
@@ -317,8 +351,12 @@ async def get_moving_average(days: int = 7, exchange: str = None):
|
||||
"""
|
||||
|
||||
if exchange:
|
||||
try:
|
||||
exchange = validate_exchange(exchange)
|
||||
except ValueError:
|
||||
raise HTTPException(status_code=400, detail="Ungueltiger Exchange-Name")
|
||||
query += f" and exchange = '{exchange}'"
|
||||
|
||||
|
||||
query += " order by date asc, exchange asc"
|
||||
|
||||
data = query_questdb(query, timeout=5)
|
||||
@@ -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]:
|
||||
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"""
|
||||
select
|
||||
@@ -423,22 +465,55 @@ async def get_analytics(
|
||||
continents: str = None
|
||||
):
|
||||
"""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"]
|
||||
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([
|
||||
group_by in ["name", "continent", "sector"] + composite_keys,
|
||||
sub_group_by in ["name", "continent", "sector"] + composite_keys,
|
||||
continents is not None
|
||||
])
|
||||
|
||||
|
||||
t_prefix = "t." if needs_metadata else ""
|
||||
m_prefix = "m." if needs_metadata else ""
|
||||
|
||||
|
||||
metrics_map = {
|
||||
"volume": f"sum({t_prefix}price * {t_prefix}quantity)",
|
||||
"count": f"count(*)",
|
||||
"avg_price": f"avg({t_prefix}price)"
|
||||
}
|
||||
|
||||
|
||||
groups_map = {
|
||||
"day": f"date_trunc('day', {t_prefix}timestamp)",
|
||||
"month": f"date_trunc('month', {t_prefix}timestamp)",
|
||||
@@ -450,35 +525,35 @@ async def get_analytics(
|
||||
"exchange_continent": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}continent, 'Unknown'))" if needs_metadata else "'Unknown'",
|
||||
"exchange_sector": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}sector, 'Unknown'))" if needs_metadata else "'Unknown'"
|
||||
}
|
||||
|
||||
|
||||
selected_metric = metrics_map.get(metric, metrics_map["volume"])
|
||||
selected_group = groups_map.get(group_by, groups_map["day"])
|
||||
|
||||
|
||||
query = f"select {selected_group} as label"
|
||||
|
||||
|
||||
if sub_group_by and sub_group_by in groups_map:
|
||||
query += f", {groups_map[sub_group_by]} as sub_label"
|
||||
|
||||
|
||||
if metric == 'all':
|
||||
query += f", count(*) as value_count, sum({t_prefix}price * {t_prefix}quantity) as value_volume from trades"
|
||||
else:
|
||||
query += f", {selected_metric} as value from trades"
|
||||
if needs_metadata:
|
||||
query += " t left join metadata m on t.isin = m.isin"
|
||||
|
||||
|
||||
query += " where 1=1"
|
||||
|
||||
|
||||
if date_from:
|
||||
query += f" and {t_prefix}timestamp >= '{date_from}'"
|
||||
if date_to:
|
||||
query += f" and {t_prefix}timestamp <= '{date_to}'"
|
||||
|
||||
if isins:
|
||||
isins_list = ",".join([f"'{i.strip()}'" for i in isins.split(",")])
|
||||
|
||||
if validated_isins:
|
||||
isins_list = ",".join([f"'{i}'" for i in validated_isins])
|
||||
query += f" and {t_prefix}isin in ({isins_list})"
|
||||
|
||||
if continents and needs_metadata:
|
||||
cont_list = ",".join([f"'{c.strip()}'" for c in continents.split(",")])
|
||||
if sanitized_continents and needs_metadata:
|
||||
cont_list = ",".join([f"'{c}'" for c in sanitized_continents])
|
||||
query += f" and {m_prefix}continent in ({cont_list})"
|
||||
|
||||
query += f" group by {selected_group}"
|
||||
@@ -493,7 +568,8 @@ async def get_analytics(
|
||||
@app.get("/api/metadata/search")
|
||||
async def search_metadata(q: str):
|
||||
"""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)
|
||||
return format_questdb_response(data)
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ import os
|
||||
import requests
|
||||
from typing import Dict, List, Tuple, Optional
|
||||
import pandas as pd
|
||||
from src.utils.validation import validate_table_name, validate_exchange
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
@@ -499,8 +500,12 @@ class AnalyticsWorker:
|
||||
"""
|
||||
|
||||
if exchange_filter:
|
||||
try:
|
||||
exchange_filter = validate_exchange(exchange_filter)
|
||||
except ValueError:
|
||||
continue
|
||||
query += f" and exchange = '{exchange_filter}'"
|
||||
|
||||
|
||||
query += f" group by date_trunc('day', timestamp), {group_by_field}"
|
||||
|
||||
data = self.query_questdb(query)
|
||||
@@ -764,6 +769,7 @@ class AnalyticsWorker:
|
||||
|
||||
def get_existing_dates(self, table_name: str) -> set:
|
||||
"""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}"
|
||||
data = self.query_questdb(query)
|
||||
if not data:
|
||||
@@ -875,12 +881,13 @@ class AnalyticsWorker:
|
||||
|
||||
for table in tables:
|
||||
try:
|
||||
table = validate_table_name(table)
|
||||
# QuestDB DELETE syntax
|
||||
delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'"
|
||||
response = requests.get(
|
||||
f"{self.questdb_url}/exec",
|
||||
f"{self.db_url}/exec",
|
||||
params={'query': delete_query},
|
||||
auth=self.auth,
|
||||
auth=DB_AUTH,
|
||||
timeout=30
|
||||
)
|
||||
if response.status_code == 200:
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
import requests
|
||||
import time
|
||||
import logging
|
||||
from typing import List
|
||||
from ..exchanges.base import Trade
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class DatabaseClient:
|
||||
def __init__(self, host: str = "localhost", port: int = 9000, user: str = None, password: str = None):
|
||||
self.host = host
|
||||
@@ -15,8 +18,8 @@ class DatabaseClient:
|
||||
return
|
||||
|
||||
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):
|
||||
batch = trades[i:i + batch_size]
|
||||
lines = []
|
||||
@@ -25,34 +28,34 @@ class DatabaseClient:
|
||||
try:
|
||||
symbol = trade.symbol.replace(" ", "\\ ").replace(",", "\\,")
|
||||
exchange = trade.exchange
|
||||
|
||||
|
||||
line = f"trades,exchange={exchange},symbol={symbol},isin={trade.isin} " \
|
||||
f"price={trade.price},quantity={trade.quantity} " \
|
||||
f"{int(trade.timestamp.timestamp() * 1e9)}"
|
||||
lines.append(line)
|
||||
except Exception as e:
|
||||
print(f"Error formating trade {trade}: {e}")
|
||||
logger.error(f"Fehler beim Formatieren von Trade {trade}: {e}")
|
||||
continue
|
||||
|
||||
if not lines:
|
||||
continue
|
||||
|
||||
payload = "\n".join(lines) + "\n"
|
||||
|
||||
|
||||
try:
|
||||
response = requests.post(
|
||||
self.url,
|
||||
data=payload,
|
||||
self.url,
|
||||
data=payload,
|
||||
params={'precision': 'ns'},
|
||||
auth=self.auth
|
||||
)
|
||||
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:
|
||||
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:
|
||||
print(f"Could not connect to QuestDB at {self.url}: {e}")
|
||||
# Fallback: print to console or save to file
|
||||
logger.error(f"Verbindung zu QuestDB fehlgeschlagen ({self.url}): {e}")
|
||||
# Fallback: in Datei speichern
|
||||
self._fallback_save(batch)
|
||||
|
||||
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 time
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import List, Optional
|
||||
from .base import BaseExchange, Trade
|
||||
import re
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Rate-Limiting Konfiguration
|
||||
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
||||
|
||||
@@ -176,9 +179,9 @@ class BoersenagBase(BaseExchange):
|
||||
|
||||
except requests.exceptions.HTTPError as e:
|
||||
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:
|
||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
||||
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||
|
||||
return trades
|
||||
|
||||
@@ -265,9 +268,9 @@ class BoersenagBase(BaseExchange):
|
||||
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}] 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
|
||||
urls = self._generate_file_urls(target_date)
|
||||
@@ -281,7 +284,7 @@ class BoersenagBase(BaseExchange):
|
||||
if trades:
|
||||
all_trades.extend(trades)
|
||||
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
|
||||
break
|
||||
|
||||
@@ -293,7 +296,7 @@ class BoersenagBase(BaseExchange):
|
||||
if i > 20 and successful == 0:
|
||||
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
|
||||
|
||||
|
||||
|
||||
@@ -3,11 +3,14 @@ import gzip
|
||||
import csv
|
||||
import io
|
||||
import time
|
||||
import logging
|
||||
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
|
||||
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}"
|
||||
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:
|
||||
print(f"[GETTEX] Error fetching page: {e}")
|
||||
logger.error(f"[GETTEX] Error fetching page: {e}")
|
||||
|
||||
return files
|
||||
|
||||
@@ -148,11 +151,11 @@ class GettexExchange(BaseExchange):
|
||||
date_str = parts[1] # YYYYMMDD
|
||||
|
||||
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
|
||||
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!
|
||||
# Format: ISIN,Zeit,Währung,Preis,Menge
|
||||
@@ -166,24 +169,24 @@ class GettexExchange(BaseExchange):
|
||||
if trade:
|
||||
trades.append(trade)
|
||||
else:
|
||||
if i < 3: # Zeige nur erste paar Fehler
|
||||
print(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
||||
if i < 3:
|
||||
logger.debug(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
||||
except Exception as e:
|
||||
parse_errors += 1
|
||||
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
|
||||
|
||||
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:
|
||||
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:
|
||||
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:
|
||||
print(f"[GETTEX] Error downloading {filename}: {e}")
|
||||
logger.error(f"[GETTEX] Error downloading {filename}: {e}")
|
||||
|
||||
return trades
|
||||
|
||||
@@ -406,9 +409,9 @@ class GettexExchange(BaseExchange):
|
||||
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}] 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
|
||||
page_files = self._get_file_list_from_page()
|
||||
@@ -434,10 +437,10 @@ class GettexExchange(BaseExchange):
|
||||
hour = int(parts[2])
|
||||
if hour < 3:
|
||||
target_files.append(f)
|
||||
except:
|
||||
except (ValueError, IndexError):
|
||||
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)
|
||||
for i, f in enumerate(target_files):
|
||||
@@ -450,9 +453,9 @@ class GettexExchange(BaseExchange):
|
||||
|
||||
# Fallback: Versuche erwartete Dateinamen
|
||||
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)
|
||||
print(f"[{self.name}] Trying {len(expected_files)} potential files")
|
||||
logger.info(f"[{self.name}] Trying {len(expected_files)} potential files")
|
||||
|
||||
successful_files = 0
|
||||
for filename in expected_files:
|
||||
@@ -461,9 +464,9 @@ class GettexExchange(BaseExchange):
|
||||
all_trades.extend(trades)
|
||||
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
|
||||
|
||||
@@ -506,12 +509,12 @@ class GettexExchange(BaseExchange):
|
||||
continue
|
||||
|
||||
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:
|
||||
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:
|
||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
||||
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||
|
||||
return trades
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
import requests
|
||||
import csv
|
||||
import io
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from typing import List
|
||||
from .base import BaseExchange, Trade
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class LSExchange(BaseExchange):
|
||||
@property
|
||||
def name(self) -> str:
|
||||
@@ -15,23 +18,23 @@ class LSExchange(BaseExchange):
|
||||
if include_yesterday:
|
||||
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lstctradesyesterday")
|
||||
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lsxtradesyesterday")
|
||||
|
||||
|
||||
headers = {
|
||||
'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Accept': 'application/json',
|
||||
'Referer': 'https://www.ls-tc.de/'
|
||||
}
|
||||
|
||||
|
||||
all_trades = []
|
||||
for url in endpoints:
|
||||
try:
|
||||
response = requests.get(url, headers=headers)
|
||||
response.raise_for_status()
|
||||
|
||||
|
||||
f = io.StringIO(response.text)
|
||||
# Header: isin;displayName;tradeTime;price;currency;size;orderId
|
||||
reader = csv.DictReader(f, delimiter=';')
|
||||
|
||||
|
||||
for item in reader:
|
||||
try:
|
||||
price = float(item['price'].replace(',', '.'))
|
||||
@@ -39,11 +42,11 @@ class LSExchange(BaseExchange):
|
||||
isin = item['isin']
|
||||
symbol = item['displayName']
|
||||
time_str = item['tradeTime']
|
||||
|
||||
|
||||
# Format: 2026-01-23T07:30:00.992000Z
|
||||
ts_str = time_str.replace('Z', '+00:00')
|
||||
timestamp = datetime.fromisoformat(ts_str)
|
||||
|
||||
|
||||
all_trades.append(Trade(
|
||||
exchange=self.name,
|
||||
symbol=symbol,
|
||||
@@ -52,8 +55,9 @@ class LSExchange(BaseExchange):
|
||||
quantity=quantity,
|
||||
timestamp=timestamp
|
||||
))
|
||||
except Exception:
|
||||
except (ValueError, KeyError) as e:
|
||||
logger.debug(f"Fehler beim Parsen einer LS-Zeile: {e}")
|
||||
continue
|
||||
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
|
||||
|
||||
@@ -3,11 +3,14 @@ import gzip
|
||||
import json
|
||||
import csv
|
||||
import io
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import List, Optional
|
||||
from .base import BaseExchange, Trade
|
||||
from bs4 import BeautifulSoup
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Browser User-Agent (Vollständiger Browser-Fingerprint für Stuttgart)
|
||||
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',
|
||||
@@ -84,7 +87,7 @@ class StuttgartExchange(BaseExchange):
|
||||
files = self._generate_expected_urls()
|
||||
|
||||
except Exception as e:
|
||||
print(f"[STU] Error fetching page: {e}")
|
||||
logger.error(f"[STU] Error fetching page: {e}")
|
||||
files = self._generate_expected_urls()
|
||||
|
||||
return files
|
||||
@@ -194,13 +197,13 @@ class StuttgartExchange(BaseExchange):
|
||||
if trade:
|
||||
trades.append(trade)
|
||||
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:
|
||||
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:
|
||||
print(f"[STU] Error downloading {url}: {e}")
|
||||
logger.error(f"[STU] Error downloading {url}: {e}")
|
||||
|
||||
return trades
|
||||
|
||||
@@ -275,7 +278,7 @@ class StuttgartExchange(BaseExchange):
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
print(f"[STU] Error parsing JSON record: {e}")
|
||||
logger.debug(f"[STU] Error parsing JSON record: {e}")
|
||||
return None
|
||||
|
||||
def _parse_csv_row(self, row: dict) -> Optional[Trade]:
|
||||
@@ -331,7 +334,7 @@ class StuttgartExchange(BaseExchange):
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
print(f"[STU] Error parsing CSV row: {e}")
|
||||
logger.debug(f"[STU] Error parsing CSV row: {e}")
|
||||
return None
|
||||
|
||||
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)
|
||||
|
||||
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
|
||||
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
|
||||
target_links = self._filter_files_for_date(all_links, target_date)
|
||||
@@ -380,7 +383,7 @@ class StuttgartExchange(BaseExchange):
|
||||
# Fallback: Versuche alle 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
|
||||
successful = 0
|
||||
@@ -389,9 +392,9 @@ class StuttgartExchange(BaseExchange):
|
||||
if trades:
|
||||
all_trades.extend(trades)
|
||||
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")
|
||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||
logger.info(f"[{self.name}] Successfully processed {successful} files")
|
||||
logger.info(f"[{self.name}] Total trades fetched: {len(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