import pandas as pd import numpy as np import yfinance as yf import os try: import talib as ta except ImportError: ta = None from datetime import datetime, timedelta from newsapi import NewsApiClient from vaderSentiment.vaderSentiment import SentimentIntensityAnalyzer from sklearn.preprocessing import MinMaxScaler from alpha_vantage.timeseries import TimeSeries import time import logging # Configure logging logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') def print_log(message, level='INFO'): if level == 'INFO': logging.info(message) elif level == 'WARNING': logging.warning(message) elif level == 'ERROR': logging.error(message) else: logging.debug(message) analyzer = SentimentIntensityAnalyzer() def load_data(data_src='yahoo', ticker='AAPL', start='2020-01-01', end='2023-01-01', interval='1d', file_upload=None, alpha_api_key=None): try: print_log(f"Loading data: source={data_src}, ticker={ticker}, start={start}, end={end}, interval={interval}, file_upload={'set' if file_upload else 'unset'}, alpha_api_key={'set' if alpha_api_key else 'unset'}") start_date = pd.to_datetime(start) end_date = pd.to_datetime(end) if start_date >= end_date: raise ValueError(f"Start date {start} must be before end date {end}") if end_date > datetime.now(): print_log(f"End date {end} is in the future. Using current date as end date.", 'WARNING') end_date = datetime.now() df = pd.DataFrame() if data_src == 'csv' and file_upload: try: file_path = getattr(file_upload, 'name', file_upload) print_log(f"Loading CSV from {file_path}") df = pd.read_csv(file_path) if 'Date' not in df.columns: raise ValueError("CSV must contain a 'Date' column") df['Date'] = pd.to_datetime(df['Date']).dt.tz_localize(None) df = df.set_index('Date') if 'Close' not in df.columns and 'value' not in df.columns: raise ValueError("CSV must contain 'Close' or 'value' column") if 'Close' in df.columns: df = df.rename(columns={'Close': 'value'}) if df.empty: raise ValueError(f"CSV data is empty for {ticker}") if df['value'].isna().all(): raise ValueError(f"CSV 'value' column contains only NaNs for {ticker}") except Exception as e: print_log(f"Failed to load CSV {file_path}: {str(e)}", 'ERROR') raise ValueError(f"Failed to load CSV: {str(e)}") else: print_log(f"Fetching data for {ticker} from Yahoo Finance") try: df = yf.download(ticker, start=start_date, end=end_date, interval=interval, progress=False, auto_adjust=False) if isinstance(df.columns, pd.MultiIndex): df.columns = df.columns.droplevel(1) if df.empty: raise ValueError(f"No data returned from Yahoo Finance for {ticker}") if 'Close' not in df.columns: raise ValueError(f"Yahoo Finance data missing 'Close' column for {ticker}") df = df.rename(columns={'Close': 'value'}) # The index is already datetime, no need to create a 'Date' column and then reset if df['value'].isna().all(): raise ValueError(f"Yahoo Finance 'value' column contains only NaNs for {ticker}") if df['value'].empty: raise ValueError(f"Yahoo Finance 'value' column is empty for {ticker}") except Exception as e: print_log(f"Yahoo Finance failed for {ticker}: {str(e)}", 'ERROR') raise ValueError(f"Yahoo Finance failed for {ticker}: {str(e)}") # Alpha Vantage data loading (if applicable) # Note: Alpha Vantage data loading logic is commented out for now to simplify debugging # if interval in ['1m', '5m', '15m', '30m', '60m'] and alpha_api_key: # print_log(f"Attempting Alpha Vantage for {ticker}, interval {interval}") # try: # ts = TimeSeries(key=alpha_api_key, output_format='pandas') # df_av, _ = ts.get_intraday(symbol=ticker, interval=interval, outputsize='full') # if df_av.empty: # raise ValueError(f"No data returned from Alpha Vantage for {ticker}") # if '4. close' not in df_av.columns: # raise ValueError(f"Alpha Vantage data missing '4. close' column for {ticker}") # df_av = df_av.rename(columns={'4. close': 'value', '1. open': 'Open', '2. high': 'High', '3. low': 'Low', '5. volume': 'Volume'}) # Standardize column names # df_av['Date'] = pd.to_datetime(df_av.index) # df_av = df_av.reset_index(drop=True) # if df_av['value'].isna().all(): # raise ValueError(f"Alpha Vantage 'value' column contains only NaNs for {ticker}") # if df_av['value'].empty: # raise ValueError(f"Alpha Vantage 'value' column is empty for {ticker}") # df = df_av # Use Alpha Vantage data if successful # except Exception as e: # print_log(f"Alpha Vantage failed: {str(e)}, using Yahoo Finance data", 'WARNING') if df.empty: raise ValueError(f"No data loaded for {ticker} from {data_src}") # Ensure index is DatetimeIndex and sorted if not isinstance(df.index, pd.DatetimeIndex): df.index = pd.to_datetime(df.index) df = df.sort_index() required_cols = ['Open', 'High', 'Low', 'value', 'Volume'] for col in required_cols: if col not in df.columns: df[col] = np.nan # Add missing columns with NaNs if 'value' not in df.columns: raise ValueError(f"Target column 'value' is missing for {ticker}") if df['value'].isna().all(): raise ValueError(f"Target column 'value' contains only NaNs for {ticker}") if df['value'].empty: raise ValueError(f"Target column 'value' is empty for {ticker}") print_log(f"Data loaded for {ticker} with date range: {df.index.min()} to {df.index.max()}, shape: {df.shape}") return df except Exception as e: print_log(f"Error in load_data for {ticker}: {str(e)}", 'ERROR') raise ValueError(f"Failed to load data for {ticker}: {str(e)}") def add_technical_indicators(df, selected_indicators): try: print_log(f"Starting add_technical_indicators with indicators: {selected_indicators}") if df.empty: print_log("DataFrame is empty, skipping technical indicator calculation.", "WARNING") return df, [] # Ensure columns are numeric and handle missing ones for col in ['Open', 'High', 'Low', 'value', 'Volume']: if col not in df.columns: df[col] = np.nan df[col] = pd.to_numeric(df[col], errors='coerce') # Drop rows with NaN in core columns after indicator calculation df.dropna(subset=['Open', 'High', 'Low', 'value', 'Volume'], inplace=True) if df.empty: print_log("DataFrame is empty after dropping NaNs for technical indicators.", "WARNING") return df, [] if ta is None: print_log("TA-Lib not available. Cannot compute indicators. Falling back to 'value'.", 'ERROR') return df, [] close = df['value'].values high = df['High'].values low = df['Low'].values volume = df['Volume'].values open_ = df['Open'].values indicator_map = { 'rsi': {'func': ta.RSI, 'inputs': ['close'], 'params': {'timeperiod': 14}, 'output': ['rsi_14']}, 'macd': {'func': ta.MACD, 'inputs': ['close'], 'params': {'fastperiod': 12, 'slowperiod': 26, 'signalperiod': 9}, 'output': ['macd_12_26_9', 'macds_12_26_9', 'macdh_12_26_9']}, 'bbands': {'func': ta.BBANDS, 'inputs': ['close'], 'params': {'timeperiod': 20, 'nbdevup': 2, 'nbdevdn': 2}, 'output': ['bbu_20_2.0', 'bbm_20_2.0', 'bbl_20_2.0']}, 'sma': {'func': ta.SMA, 'inputs': ['close'], 'params': {'timeperiod': 20}, 'output': ['sma_20']}, 'ema': {'func': ta.EMA, 'inputs': ['close'], 'params': {'timeperiod': 20}, 'output': ['ema_20']}, 'atr': {'func': ta.ATR, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 14}, 'output': ['atr_14']}, 'stoch': {'func': ta.STOCH, 'inputs': ['high', 'low', 'close'], 'params': {'fastk_period': 14, 'slowk_period': 3, 'slowd_period': 3}, 'output': ['stochk_14_3_3', 'stochd_14_3_3']}, 'adx': {'func': ta.ADX, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 14}, 'output': ['adx_14']}, 'willr': {'func': ta.WILLR, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 14}, 'output': ['willr_14']}, 'cci': {'func': ta.CCI, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 20}, 'output': ['cci_20']}, 'pdi': {'func': ta.PLUS_DI, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 14}, 'output': ['pdi_14']}, 'mdi': {'func': ta.MINUS_DI, 'inputs': ['high', 'low', 'close'], 'params': {'timeperiod': 14}, 'output': ['mdi_14']} } input_dict = {'close': close, 'high': high, 'low': low, 'open': open_, 'volume': volume} valid_indicators = [] for ind in selected_indicators: if ind in indicator_map: print_log(f"Computing indicator: {ind}") config = indicator_map[ind] func = config['func'] inputs = config['inputs'] params = config['params'] try: input_arrays = [input_dict[inp] for inp in inputs] result = func(*input_arrays, **params) if isinstance(result, tuple): for j, (res, out_col) in enumerate(zip(result, config['output'])): if isinstance(res, np.ndarray) and len(res) == len(df): df[out_col] = res nan_count = np.isnan(res).sum() if nan_count < len(res) * 0.5: valid_indicators.append(out_col) else: print_log(f"{out_col} has excessive NaNs: {nan_count}/{len(res)}. Excluding from valid indicators.", 'WARNING') else: print_log(f"Invalid output for {out_col}: {type(res)}, length: {len(res) if hasattr(res, '__len__') else 'N/A'}", 'WARNING') else: if isinstance(result, np.ndarray) and len(result) == len(df): df[config['output'][0]] = result nan_count = np.isnan(result).sum() if nan_count < len(result) * 0.5: valid_indicators.append(config['output'][0]) else: print_log(f"{config['output'][0]} has excessive NaNs: {nan_count}/{len(result)}. Excluding from valid indicators.", 'WARNING') else: print_log(f"Invalid output for {ind}: {type(result)}, length: {len(result) if hasattr(result, '__len__') else 'N/A'}", 'WARNING') except Exception as e: print_log(f"Error computing {ind}: {str(e)}", 'ERROR') else: print_log(f"Indicator {ind} not supported by TA-Lib", 'WARNING') # Drop rows with NaN in 'value', preserve valid indicators initial_rows = len(df) df = df.dropna(subset=['value']).reset_index(drop=False) # Keep index as a column for now print_log(f"Dropped {initial_rows - len(df)} rows with NaN in 'value'") # Drop columns with excessive NaNs, but protect 'value' for col in df.columns: if col not in ['Date', 'Open', 'High', 'Low', 'value', 'Volume']: nan_ratio = df[col].isna().mean() if nan_ratio > 0.5: print_log(f"Dropping {col} due to excessive NaNs: {nan_ratio:.2%}", 'WARNING') df = df.drop(columns=[col]) if col in valid_indicators: valid_indicators.remove(col) if not valid_indicators: print_log("No valid indicators computed. Falling back to 'value'.", "WARNING") valid_indicators.append('value') print_log(f"Valid indicators: {valid_indicators}") print_log(f"Technical indicators added successfully, shape: {df.shape}") df.set_index('Date', inplace=True) # Set index back to Date after all processing return df, valid_indicators except Exception as e: print_log(f"Error in add_technical_indicators: {str(e)}", 'ERROR') # If an error occurs, return the original DataFrame to prevent further errors return df, [] def add_sentiment(df, ticker, news_api_key, start_date, end_date): try: print_log(f"Starting add_sentiment for {ticker} from {start_date} to {end_date}") if not news_api_key: print_log("News API key not provided. Skipping sentiment analysis.", "WARNING") df['sentiment_score'] = 0.0 return df newsapi = NewsApiClient(api_key=news_api_key) all_articles = [] current_date = pd.to_datetime(start_date) end_date = pd.to_datetime(end_date) while current_date <= end_date: from_param = current_date.strftime("%Y-%m-%d") to_param = (current_date + timedelta(days=1)).strftime("%Y-%m-%d") print_log(f"Fetching news for {ticker} from {from_param} to {to_param}") try: articles = newsapi.get_everything(q=ticker, language='en', sort_by='relevancy', from_param=from_param, to=to_param) all_articles.extend(articles["articles"]) except Exception as e: print_log(f"Error fetching news for {ticker} on {from_param}: {str(e)}", 'ERROR') current_date += timedelta(days=1) time.sleep(0.1) if not all_articles: print_log(f"No articles found for {ticker}. Setting sentiment to 0.", 'WARNING') df['sentiment_score'] = 0.0 return df sentiment_data = [] for article in all_articles: if article["publishedAt"] and article["description"]: date = pd.to_datetime(article["publishedAt"]).tz_localize(None).date() text = article["description"] vs = analyzer.polarity_scores(text) sentiment_data.append({"Date": date, "sentiment_score": vs["compound"]}) sentiment_df = pd.DataFrame(sentiment_data) sentiment_df["Date"] = pd.to_datetime(sentiment_df["Date"]) sentiment_df = sentiment_df.groupby("Date")["sentiment_score"].mean().reset_index() df.reset_index(inplace=True) df['Date'] = pd.to_datetime(df['Date']) df = pd.merge(df, sentiment_df, on="Date", how="left") df['sentiment_score'] = df['sentiment_score'].fillna(0.0) df.set_index('Date', inplace=True) print_log(f"Sentiment analysis completed for {ticker}. Added sentiment_score column.") return df except Exception as e: print_log(f"Error in add_sentiment for {ticker}: {str(e)}", 'ERROR') df['sentiment_score'] = 0.0 return df def preprocess_data(df, features, target, window_size, horizon): try: print_log(f"Starting preprocessing: features={features}, target={target}, window={window_size}, horizon={horizon}") # Ensure the DataFrame index is a DatetimeIndex if not isinstance(df.index, pd.DatetimeIndex): raise ValueError("DataFrame index must be a DatetimeIndex for preprocessing.") # Filter features to only include those present in the DataFrame columns updated_feature_cols = [f for f in features if f in df.columns] if not updated_feature_cols: raise ValueError("No valid features found in DataFrame after indicator calculation.") full_features = updated_feature_cols + [target] data = df[full_features].copy() data.dropna(inplace=True) if data.empty: raise ValueError("DataFrame is empty after dropping NaNs. Cannot proceed with scaling.") feature_scaler = MinMaxScaler() target_scaler = MinMaxScaler() data_features_scaled = feature_scaler.fit_transform(data[updated_feature_cols]) data_target_scaled = target_scaler.fit_transform(data[[target]]) full_scaled = np.hstack((data_features_scaled, data_target_scaled)) target_idx = len(updated_feature_cols) X, y = [], [] for i in range(len(full_scaled) - window_size - horizon + 1): X.append(full_scaled[i:i + window_size]) y.append(full_scaled[i + window_size:i + window_size + horizon, target_idx]) X = np.array(X) y = np.array(y) if X.shape[0] == 0 or y.shape[0] == 0: raise ValueError(f"Insufficient data after preprocessing. Data length: {len(full_scaled)}, window_size: {window_size}, horizon: {horizon}") print_log(f"Preprocessed data: X.shape={X.shape}, y.shape={y.shape}, Final features: {full_features}, Target idx: {target_idx}") return X, y, feature_scaler, target_scaler, full_features, target_idx, None, updated_feature_cols except Exception as e: print_log(f"Preprocessing error: {str(e)}", 'ERROR') raise ValueError(f"Preprocessing failed: {str(e)}")