Files
hyper/resampler.py
DiTus 1a95fe1caa Clean up unused files, organize structure, update docs
- Delete obsolete files: data_fetcher_old.py, market_old.py, base_strategy.py (root),
  strategy_sma_cross.py, and old architecture remnants (address_monitor.py,
  position_monitor.py, trade_log.py, wallet_data.py, whale_tracker.py)
- Delete zero-byte Docker artifacts and runtime files (clp_hedger.log,
  clp_hedger/hedge_status.json)
- Move one-off utility scripts to scripts/ directory
- Move example/template files to .temp/ directory
- Update .gitignore: add entries for clp_hedger.log, clp_hedger/hedge_status.json,
  Docker layer hash files, Using, Running, and backups/
- Update .dockerignore: add clp_hedger.log, clp_hedger/hedge_status.json, backups/
- Create example config files: _data/strategies.json.example,
  _data/backtesting_conf.json.example, _data/coin_precision.json.example
- Update GEMINI.md: remove outdated session summaries and duplicate review section
- Update review.md: add cleanup status section, update remaining recommendations
- Update MIGRATION_PLAN.md: mark completed phases, update file references
- Update DOCKER_MIGRATION_GUIDE.md: update import_csv.py path reference
2026-08-05 09:50:36 +02:00

242 lines
11 KiB
Python

import argparse
import logging
import os
import sys
import warnings
import db
import pandas as pd
import json
from datetime import datetime, timezone, timedelta
warnings.filterwarnings("ignore", message="pandas only supports SQLAlchemy")
# Assuming logging_utils.py is in the same directory
from logging_utils import setup_logging
class Resampler:
"""
Reads new 1-minute candle data from the PostgreSQL database, resamples it to
various timeframes, and upserts the new candles to the corresponding tables,
preventing data duplication.
"""
def __init__(self, log_level: str, coins: list, timeframes: dict):
setup_logging(log_level, 'Resampler')
self.db_path = os.environ.get("PG_CONN_STR", "postgresql://hyper:hyper@localhost:5432/hyper")
self.status_file_path = os.path.join("_data", "resampling_status.json")
self.coins_to_process = coins
self.timeframes = timeframes
self.aggregation_logic = {
'open': 'first',
'high': 'max',
'low': 'min',
'close': 'last',
'volume': 'sum',
'number_of_trades': 'sum'
}
self.resampling_status = self._load_existing_status()
self.job_start_time = None
self._ensure_tables_exist()
def _ensure_tables_exist(self):
"""
Ensures all resampled tables exist with the correct schema.
Uses db.create_candle_table() which is idempotent.
"""
conn = db.get_connection()
for coin in self.coins_to_process:
for tf_name in self.timeframes.keys():
table_name = db.sanitize_table_name(coin, tf_name)
db.create_candle_table(conn, table_name)
conn.close()
logging.info("All resampled table schemas verified.")
def _load_existing_status(self) -> dict:
"""Loads the existing status file if it exists, otherwise returns an empty dict."""
if os.path.exists(self.status_file_path):
try:
with open(self.status_file_path, 'r', encoding='utf-8') as f:
logging.debug(f"Loading existing status from '{self.status_file_path}'")
return json.load(f)
except (IOError, json.JSONDecodeError) as e:
logging.warning(f"Could not read existing status file. Starting fresh. Error: {e}")
return {}
def run(self):
"""
Main execution function to process all configured coins and update the database.
"""
self.job_start_time = datetime.now(timezone.utc)
logging.info(f"--- Resampling job started at {self.job_start_time.strftime('%Y-%m-%d %H:%M:%S %Z')} ---")
if '1m' in self.timeframes:
logging.debug("Ignoring '1m' timeframe as it is the source resolution.")
del self.timeframes['1m']
if not self.timeframes:
logging.warning("No timeframes to process after filtering. Exiting job.")
return
conn = db.get_connection()
try:
logging.debug(f"Processing {len(self.coins_to_process)} coins...")
for coin in self.coins_to_process:
logging.info(f"--- Processing {coin} ---")
try:
for tf_name, tf_code in self.timeframes.items():
target_table_name = db.sanitize_table_name(coin, tf_name)
source_table_name = db.sanitize_table_name(coin, "1m")
logging.info(f" Resampling {coin} -> {tf_name}")
last_timestamp_ms = self._get_last_timestamp(conn, target_table_name)
query = f'SELECT * FROM "{source_table_name}"'
params = ()
if last_timestamp_ms:
query += ' WHERE timestamp_ms >= %s'
# Go back one interval to rebuild the last (potentially partial) candle
try:
interval_delta_ms = pd.to_timedelta(tf_code).total_seconds() * 1000
except ValueError:
# Fall back to a safe 32-day lookback for special timeframes
interval_delta_ms = timedelta(days=32).total_seconds() * 1000
query_start_ms = last_timestamp_ms - interval_delta_ms
params = (query_start_ms,)
df_1m = pd.read_sql(query, conn, params=params, parse_dates=['datetime_utc'])
if df_1m.empty:
logging.debug(f" -> No new 1-minute data for {tf_name}. Table is up to date.")
continue
df_1m.set_index('datetime_utc', inplace=True)
resampled_df = df_1m.resample(tf_code).agg(self.aggregation_logic)
resampled_df.dropna(how='all', inplace=True)
if not resampled_df.empty:
records_to_upsert = []
for index, row in resampled_df.iterrows():
records_to_upsert.append((
index.strftime('%Y-%m-%d %H:%M:%S'),
int(index.timestamp() * 1000), # Generate timestamp_ms
float(row['open']), float(row['high']), float(row['low']), float(row['close']),
float(row['volume']), int(row['number_of_trades'])
))
db.upsert_candles(conn, target_table_name, records_to_upsert)
logging.debug(f" -> Upserted {len(resampled_df)} candles into '{target_table_name}'.")
if coin not in self.resampling_status: self.resampling_status[coin] = {}
total_candles = int(self._get_table_count(conn, target_table_name))
self.resampling_status[coin][tf_name] = {
"last_candle_utc": resampled_df.index[-1].strftime('%Y-%m-%d %H:%M:%S'),
"total_candles": total_candles
}
except Exception as e:
logging.error(f"Failed to process coin '{coin}': {e}")
finally:
conn.close()
self._log_summary()
self._save_status()
logging.info(f"--- Resampling job finished at {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M:%S %Z')} ---")
def _log_summary(self):
"""Logs a summary of the total candles for each timeframe."""
logging.info("--- Resampling Job Summary ---")
timeframe_totals = {}
for coin, tfs in self.resampling_status.items():
if not isinstance(tfs, dict): continue
for tf_name, tf_data in tfs.items():
total = tf_data.get("total_candles", 0)
if tf_name not in timeframe_totals:
timeframe_totals[tf_name] = 0
timeframe_totals[tf_name] += total
if not timeframe_totals:
logging.info("No candles were resampled in this run.")
return
logging.info("Total candles per timeframe across all processed coins:")
for tf_name, total in sorted(timeframe_totals.items()):
logging.info(f" - {tf_name:<10}: {total:,} candles")
def _get_last_timestamp(self, conn, table_name):
"""Gets the millisecond timestamp of the last entry in a table."""
try:
# --- FIX: Query for the integer timestamp_ms, not the text datetime_utc ---
timestamp_ms = pd.read_sql(f'SELECT MAX(timestamp_ms) FROM "{table_name}"', conn).iloc[0, 0]
return int(timestamp_ms) if pd.notna(timestamp_ms) else None
except (pd.io.sql.DatabaseError, IndexError):
return None
def _get_table_count(self, conn, table_name):
"""Gets the total row count of a table."""
try:
return pd.read_sql(f'SELECT COUNT(*) FROM "{table_name}"', conn).iloc[0, 0]
except (pd.io.sql.DatabaseError, IndexError):
return 0
def _save_status(self):
"""Saves the final resampling status to a JSON file."""
if not self.resampling_status:
return
stop_time = datetime.now(timezone.utc)
self.resampling_status['job_start_time_utc'] = self.job_start_time.strftime('%Y-%m-%d %H:%M:%S')
self.resampling_status['job_stop_time_utc'] = stop_time.strftime('%Y-%m-%d %H:%M:%S')
self.resampling_status.pop('last_completed_utc', None)
try:
with open(self.status_file_path, 'w', encoding='utf-8') as f:
json.dump(self.resampling_status, f, indent=4, sort_keys=True)
logging.info(f"Successfully saved resampling status to '{self.status_file_path}'")
except IOError as e:
logging.error(f"Failed to write resampling status file: {e}")
def parse_timeframes(tf_strings: list) -> dict:
"""Converts a list of timeframe strings into a dictionary for pandas."""
tf_map = {}
for tf_str in tf_strings:
numeric_part = ''.join(filter(str.isdigit, tf_str))
unit = ''.join(filter(str.isalpha, tf_str)) # Keep case for 'M'
key = tf_str
code = ''
if unit == 'm':
code = f"{numeric_part}min"
elif unit.lower() == 'w':
code = f"{numeric_part}W-MON"
elif unit == 'M':
code = f"{numeric_part}MS"
key = f"{numeric_part}month"
elif unit.lower() in ['h', 'd']:
code = f"{numeric_part}{unit.lower()}"
else:
code = tf_str
logging.warning(f"Unrecognized timeframe unit in '{tf_str}'. Using as-is.")
tf_map[key] = code
return tf_map
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="Resample 1-minute candle data from PostgreSQL to other timeframes.")
parser.add_argument("--coins", nargs='+', required=True, help="List of coins to process.")
parser.add_argument("--timeframes", nargs='+', required=True, help="List of timeframes to generate.")
parser.add_argument("--log-level", default="normal", choices=['off', 'normal', 'debug'])
args = parser.parse_args()
timeframes_dict = parse_timeframes(args.timeframes)
resampler = Resampler(log_level=args.log_level, coins=args.coins, timeframes=timeframes_dict)
resampler.run()