From Jupyter Notebook to Live Exchange: A Production Journey
Transform your Jupyter research notebooks into production-ready trading systems deployed on Indian exchanges. Complete guide with code, infrastructure, and monitoring.
Your Jupyter notebook shows a Sharpe ratio of 2.5, consistent profits across 5 years of backtests, and beautiful equity curves. But it’s all research code: hardcoded paths, manual cell execution, no error handling, and zero monitoring. Getting from here to live trading on NSE/BSE is a journey most quantitative traders underestimate.
This comprehensive guide walks you through transforming research code into production-ready trading systems deployed on Indian exchanges, covering code refactoring, infrastructure, monitoring, and SEBI compliance.
Part 1: The Research-to-Production Gap
Typical Jupyter Notebook Structure
# Cell 1: Imports (scattered across notebook)
import pandas as pd
import numpy as np
from datetime import datetime
# Cell 2: Load data (hardcoded paths)
df = pd.read_csv('/Users/trader/Desktop/nifty_data.csv')
# Cell 3: Quick data exploration
df.head()
df.describe()
# Cell 4: Calculate indicators (no functions)
df['sma_20'] = df['Close'].rolling(20).mean()
df['sma_50'] = df['Close'].rolling(50).mean()
# Cell 5: Generate signals (messy logic)
df['signal'] = 0
for i in range(len(df)):
if df['sma_20'].iloc[i] > df['sma_50'].iloc[i]:
df.loc[i, 'signal'] = 1
else:
df.loc[i, 'signal'] = -1
# Cell 6: Backtest (no position sizing, no costs)
df['returns'] = df['Close'].pct_change()
df['strategy_returns'] = df['signal'].shift(1) * df['returns']
# Cell 7: Plot results
df['strategy_returns'].cumsum().plot()
# Cell 8: Print metrics
print(f"Total return: {df['strategy_returns'].sum()}")
Production Requirements
Code Quality:
- Modular, reusable functions
- Proper class structure
- Type hints and documentation
- Comprehensive error handling
- Unit tests and integration tests
Infrastructure:
- Automated data pipelines
- Reliable broker API integration
- Database for state management
- Logging and monitoring
- Alert systems for failures
Operational:
- Configuration management
- Secrets handling (API keys, passwords)
- Process management (restart on failure)
- Backup and recovery
- Performance monitoring
Compliance:
- SEBI algo trading requirements
- Audit trail for all trades
- Risk limit enforcement
- Kill switch implementation
Part 2: Code Refactoring
Step 1: Extract Configuration
# config.py
from dataclasses import dataclass
from pathlib import Path
import yaml
@dataclass
class BrokerConfig:
"""Broker API configuration"""
api_key: str
api_secret: str
user_id: str
password: str
totp_secret: str
@dataclass
class StrategyConfig:
"""Strategy parameters"""
symbol: str
fast_period: int
slow_period: int
position_size_pct: float
max_positions: int
stop_loss_pct: float
take_profit_pct: float
@dataclass
class SystemConfig:
"""System configuration"""
data_dir: Path
log_dir: Path
db_path: Path
run_mode: str # 'backtest', 'paper', 'live'
def load_config(config_file: str = 'config.yaml') -> dict:
"""Load configuration from YAML file"""
with open(config_file, 'r') as f:
config_dict = yaml.safe_load(f)
return {
'broker': BrokerConfig(**config_dict['broker']),
'strategy': StrategyConfig(**config_dict['strategy']),
'system': SystemConfig(**config_dict['system'])
}
# config.yaml
broker:
api_key: ${ZERODHA_API_KEY} # From environment variable
api_secret: ${ZERODHA_API_SECRET}
user_id: ${ZERODHA_USER_ID}
password: ${ZERODHA_PASSWORD}
totp_secret: ${ZERODHA_TOTP_SECRET}
strategy:
symbol: "RELIANCE"
fast_period: 20
slow_period: 50
position_size_pct: 0.02 # 2% per trade
max_positions: 5
stop_loss_pct: 0.015 # 1.5%
take_profit_pct: 0.03 # 3%
system:
data_dir: "/var/data/trading"
log_dir: "/var/log/trading"
db_path: "/var/data/trading/state.db"
run_mode: "backtest"
Step 2: Create Strategy Class
# strategy.py
from abc import ABC, abstractmethod
from typing import Dict, List, Optional
import pandas as pd
import logging
class BaseStrategy(ABC):
"""
Abstract base class for trading strategies
All strategies inherit from this
"""
def __init__(self, config: StrategyConfig):
self.config = config
self.logger = logging.getLogger(self.__class__.__name__)
self.data: Optional[pd.DataFrame] = None
self.positions: Dict[str, int] = {}
@abstractmethod
def calculate_indicators(self, data: pd.DataFrame) -> pd.DataFrame:
"""Calculate technical indicators"""
pass
@abstractmethod
def generate_signals(self, data: pd.DataFrame) -> pd.DataFrame:
"""Generate trading signals"""
pass
def validate_data(self, data: pd.DataFrame) -> bool:
"""Validate incoming data"""
required_columns = ['Open', 'High', 'Low', 'Close', 'Volume']
if not all(col in data.columns for col in required_columns):
self.logger.error(f"Missing required columns")
return False
if data.isnull().any().any():
self.logger.warning("Data contains null values")
# Handle nulls appropriately
return True
def update(self, new_data: pd.DataFrame) -> List[Dict]:
"""
Main update method called on each new data point
Returns list of signals
"""
if not self.validate_data(new_data):
return []
self.data = new_data
try:
# Calculate indicators
self.data = self.calculate_indicators(self.data)
# Generate signals
self.data = self.generate_signals(self.data)
# Extract signals for execution
signals = self._extract_signals()
return signals
except Exception as e:
self.logger.error(f"Error in strategy update: {e}", exc_info=True)
return []
def _extract_signals(self) -> List[Dict]:
"""Extract actionable signals from data"""
signals = []
latest = self.data.iloc[-1]
if latest.get('Signal') == 'BUY':
signals.append({
'action': 'BUY',
'symbol': self.config.symbol,
'quantity': self._calculate_position_size(),
'order_type': 'MARKET',
'reason': 'Strategy signal'
})
elif latest.get('Signal') == 'SELL':
signals.append({
'action': 'SELL',
'symbol': self.config.symbol,
'quantity': self.positions.get(self.config.symbol, 0),
'order_type': 'MARKET',
'reason': 'Strategy signal'
})
return signals
def _calculate_position_size(self) -> int:
"""Calculate position size based on risk parameters"""
# Implement position sizing logic
# Consider available capital, risk per trade, etc.
return 10 # Placeholder
class MovingAverageCrossover(BaseStrategy):
"""
Production version of SMA crossover strategy
"""
def calculate_indicators(self, data: pd.DataFrame) -> pd.DataFrame:
"""Calculate moving averages"""
data['SMA_Fast'] = data['Close'].rolling(
window=self.config.fast_period
).mean()
data['SMA_Slow'] = data['Close'].rolling(
window=self.config.slow_period
).mean()
return data
def generate_signals(self, data: pd.DataFrame) -> pd.DataFrame:
"""Generate crossover signals"""
data['Signal'] = 'HOLD'
# Golden cross: fast MA crosses above slow MA
golden_cross = (
(data['SMA_Fast'] > data['SMA_Slow']) &
(data['SMA_Fast'].shift(1) <= data['SMA_Slow'].shift(1))
)
# Death cross: fast MA crosses below slow MA
death_cross = (
(data['SMA_Fast'] < data['SMA_Slow']) &
(data['SMA_Fast'].shift(1) >= data['SMA_Slow'].shift(1))
)
data.loc[golden_cross, 'Signal'] = 'BUY'
data.loc[death_cross, 'Signal'] = 'SELL'
return data
Step 3: Build Data Pipeline
# data_pipeline.py
from typing import Optional
import pandas as pd
from datetime import datetime, timedelta
import yfinance as yf
import sqlite3
import logging
class DataPipeline:
"""
Robust data pipeline for production trading
"""
def __init__(self, db_path: str):
self.db_path = db_path
self.logger = logging.getLogger(__name__)
self._init_database()
def _init_database(self):
"""Initialize SQLite database for data storage"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS market_data (
symbol TEXT,
timestamp DATETIME,
open REAL,
high REAL,
low REAL,
close REAL,
volume INTEGER,
PRIMARY KEY (symbol, timestamp)
)
""")
cursor.execute("""
CREATE INDEX IF NOT EXISTS idx_symbol_timestamp
ON market_data(symbol, timestamp)
""")
conn.commit()
conn.close()
def fetch_historical_data(
self,
symbol: str,
start_date: datetime,
end_date: Optional[datetime] = None,
interval: str = 'day'
) -> pd.DataFrame:
"""
Fetch historical data with caching
"""
if end_date is None:
end_date = datetime.now()
# Check if data exists in database
cached_data = self._load_from_database(symbol, start_date, end_date)
if cached_data is not None and len(cached_data) > 0:
self.logger.info(f"Loaded {len(cached_data)} rows from cache for {symbol}")
# Check if we need to fetch new data
latest_cached = cached_data.index.max()
if latest_cached < end_date:
# Fetch missing data
new_data = self._fetch_from_source(
symbol,
start_date=latest_cached + timedelta(days=1),
end_date=end_date
)
if new_data is not None and len(new_data) > 0:
# Store new data
self._store_to_database(symbol, new_data)
# Combine with cached data
cached_data = pd.concat([cached_data, new_data])
return cached_data
else:
# Fetch all data from source
data = self._fetch_from_source(symbol, start_date, end_date)
if data is not None and len(data) > 0:
# Store to database
self._store_to_database(symbol, data)
return data
def _fetch_from_source(
self,
symbol: str,
start_date: datetime,
end_date: datetime
) -> Optional[pd.DataFrame]:
"""
Fetch data from external source (Yahoo Finance, broker API, etc.)
"""
try:
ticker = f"{symbol}.NS" # NSE symbol
data = yf.download(
ticker,
start=start_date,
end=end_date,
progress=False
)
if data.empty:
self.logger.warning(f"No data fetched for {symbol}")
return None
# Standardize column names
data.columns = ['Open', 'High', 'Low', 'Close', 'Adj Close', 'Volume']
data = data[['Open', 'High', 'Low', 'Close', 'Volume']]
self.logger.info(f"Fetched {len(data)} rows for {symbol}")
return data
except Exception as e:
self.logger.error(f"Failed to fetch data for {symbol}: {e}")
return None
def _load_from_database(
self,
symbol: str,
start_date: datetime,
end_date: datetime
) -> Optional[pd.DataFrame]:
"""Load data from SQLite database"""
try:
conn = sqlite3.connect(self.db_path)
query = """
SELECT timestamp, open, high, low, close, volume
FROM market_data
WHERE symbol = ? AND timestamp BETWEEN ? AND ?
ORDER BY timestamp
"""
data = pd.read_sql_query(
query,
conn,
params=(symbol, start_date, end_date),
index_col='timestamp',
parse_dates=['timestamp']
)
conn.close()
if data.empty:
return None
# Standardize column names
data.columns = ['Open', 'High', 'Low', 'Close', 'Volume']
return data
except Exception as e:
self.logger.error(f"Failed to load from database: {e}")
return None
def _store_to_database(self, symbol: str, data: pd.DataFrame):
"""Store data to SQLite database"""
try:
conn = sqlite3.connect(self.db_path)
# Prepare data for insertion
data_to_store = data.copy()
data_to_store['symbol'] = symbol
data_to_store['timestamp'] = data_to_store.index
data_to_store.columns = ['open', 'high', 'low', 'close', 'volume', 'symbol', 'timestamp']
# Insert with conflict handling (ignore duplicates)
data_to_store.to_sql(
'market_data',
conn,
if_exists='append',
index=False
)
conn.commit()
conn.close()
self.logger.info(f"Stored {len(data)} rows for {symbol}")
except Exception as e:
self.logger.error(f"Failed to store to database: {e}")
Step 4: Implement Execution Engine
# execution_engine.py
from typing import Dict, List, Optional
from kiteconnect import KiteConnect
import logging
from datetime import datetime
import json
class ExecutionEngine:
"""
Order execution engine with retry logic and safety checks
"""
def __init__(self, broker_config: BrokerConfig, risk_manager):
self.config = broker_config
self.risk_manager = risk_manager
self.logger = logging.getLogger(__name__)
self.kite = self._initialize_broker()
self.pending_orders = {}
def _initialize_broker(self) -> KiteConnect:
"""Initialize and authenticate with broker"""
# ... (authentication logic from previous blog posts)
pass
def execute_signals(self, signals: List[Dict]) -> List[Dict]:
"""
Execute list of trading signals with safety checks
"""
executed_orders = []
for signal in signals:
try:
# Pre-trade risk checks
if not self.risk_manager.can_trade(signal):
self.logger.warning(f"Risk check failed for signal: {signal}")
continue
# Execute order
order_id = self._place_order(signal)
if order_id:
executed_orders.append({
'signal': signal,
'order_id': order_id,
'timestamp': datetime.now()
})
# Log to audit trail
self._log_trade(signal, order_id)
except Exception as e:
self.logger.error(f"Failed to execute signal: {e}", exc_info=True)
return executed_orders
def _place_order(self, signal: Dict) -> Optional[str]:
"""Place order with error handling and retry logic"""
max_retries = 3
for attempt in range(max_retries):
try:
order_id = self.kite.place_order(
variety=self.kite.VARIETY_REGULAR,
exchange=self.kite.EXCHANGE_NSE,
tradingsymbol=signal['symbol'],
transaction_type=signal['action'], # BUY or SELL
quantity=signal['quantity'],
product=self.kite.PRODUCT_MIS, # Or CNC based on strategy
order_type=signal['order_type']
)
self.logger.info(f"Order placed successfully: {order_id}")
return order_id
except Exception as e:
self.logger.warning(f"Order placement attempt {attempt + 1} failed: {e}")
if attempt < max_retries - 1:
time.sleep(1) # Wait before retry
else:
self.logger.error(f"Order placement failed after {max_retries} attempts")
return None
def _log_trade(self, signal: Dict, order_id: str):
"""Log trade to audit trail (required for SEBI compliance)"""
audit_entry = {
'timestamp': datetime.now().isoformat(),
'signal': signal,
'order_id': order_id,
'strategy': signal.get('strategy', 'unknown')
}
# Write to audit log file
with open('/var/log/trading/audit.jsonl', 'a') as f:
f.write(json.dumps(audit_entry) + '\n')
Part 3: Production Infrastructure
Systemd Service Configuration
# /etc/systemd/system/trading-bot.service
[Unit]
Description=Automated Trading Bot
After=network.target
[Service]
Type=simple
User=trader
Group=trader
WorkingDirectory=/opt/trading
Environment="PATH=/opt/trading/venv/bin"
ExecStart=/opt/trading/venv/bin/python /opt/trading/main.py
Restart=always
RestartSec=10
StandardOutput=append:/var/log/trading/bot.log
StandardError=append:/var/log/trading/bot-error.log
# Security hardening
NoNewPrivileges=true
PrivateTmp=true
ProtectSystem=strict
ProtectHome=true
ReadWritePaths=/var/log/trading /var/data/trading
[Install]
WantedBy=multi-user.target
Main Application Entry Point
# main.py
import logging
import signal
import sys
from datetime import datetime
import time
class TradingApplication:
"""
Main trading application
"""
def __init__(self, config_file: str = 'config.yaml'):
self.config = load_config(config_file)
self.setup_logging()
self.running = False
# Initialize components
self.data_pipeline = DataPipeline(self.config['system'].db_path)
self.strategy = MovingAverageCrossover(self.config['strategy'])
self.risk_manager = RiskManager(self.config)
self.execution_engine = ExecutionEngine(
self.config['broker'],
self.risk_manager
)
# Setup graceful shutdown
signal.signal(signal.SIGTERM, self.shutdown)
signal.signal(signal.SIGINT, self.shutdown)
def setup_logging(self):
"""Configure logging"""
log_format = '%(asctime)s - %(name)s - %(levelname)s - %(message)s'
logging.basicConfig(
level=logging.INFO,
format=log_format,
handlers=[
logging.FileHandler('/var/log/trading/app.log'),
logging.StreamHandler(sys.stdout)
]
)
self.logger = logging.getLogger(__name__)
def run(self):
"""Main application loop"""
self.logger.info("Starting trading application")
self.running = True
while self.running:
try:
# Check if market is open
if not self.is_market_open():
self.logger.info("Market closed, waiting...")
time.sleep(300) # Check every 5 minutes
continue
# Fetch latest data
data = self.data_pipeline.fetch_historical_data(
symbol=self.config['strategy'].symbol,
start_date=datetime.now() - timedelta(days=100),
end_date=datetime.now()
)
# Update strategy
signals = self.strategy.update(data)
# Execute signals
if signals:
self.execution_engine.execute_signals(signals)
# Wait before next iteration
time.sleep(60) # Run every minute
except Exception as e:
self.logger.error(f"Error in main loop: {e}", exc_info=True)
time.sleep(60)
def shutdown(self, signum, frame):
"""Graceful shutdown"""
self.logger.info("Shutdown signal received")
self.running = False
# Close all positions if configured
if self.config.get('close_on_shutdown', False):
self.risk_manager.close_all_positions()
sys.exit(0)
def is_market_open(self) -> bool:
"""Check if market is open"""
# ... (market hours logic)
pass
if __name__ == '__main__':
app = TradingApplication()
app.run()
Deployment Script
#!/bin/bash
# deploy.sh
set -e
echo "Deploying trading bot to production..."
# Variables
DEPLOY_DIR="/opt/trading"
VENV_DIR="$DEPLOY_DIR/venv"
REPO_URL="[email protected]:yourcompany/trading-bot.git"
# Stop existing service
echo "Stopping existing service..."
sudo systemctl stop trading-bot || true
# Pull latest code
echo "Pulling latest code..."
cd $DEPLOY_DIR
git pull origin main
# Update dependencies
echo "Installing dependencies..."
$VENV_DIR/bin/pip install -r requirements.txt
# Run tests
echo "Running tests..."
$VENV_DIR/bin/pytest tests/
# Database migrations (if any)
echo "Running migrations..."
$VENV_DIR/bin/python migrate.py
# Start service
echo "Starting service..."
sudo systemctl start trading-bot
sudo systemctl status trading-bot
echo "Deployment complete!"
Part 4: Monitoring and Alerts
Prometheus Metrics
# metrics.py
from prometheus_client import Counter, Gauge, Histogram, start_http_server
class MetricsCollector:
"""
Collect and export metrics for monitoring
"""
def __init__(self):
# Counters
self.trades_total = Counter(
'trades_total',
'Total number of trades executed',
['action', 'symbol']
)
self.errors_total = Counter(
'errors_total',
'Total number of errors',
['error_type']
)
# Gauges
self.positions_count = Gauge(
'positions_count',
'Current number of open positions'
)
self.account_balance = Gauge(
'account_balance',
'Current account balance'
)
self.daily_pnl = Gauge(
'daily_pnl',
'Daily profit/loss'
)
# Histograms
self.order_execution_time = Histogram(
'order_execution_seconds',
'Time taken to execute orders'
)
# Start metrics server
start_http_server(9090)
def record_trade(self, action: str, symbol: str):
"""Record a trade"""
self.trades_total.labels(action=action, symbol=symbol).inc()
def record_error(self, error_type: str):
"""Record an error"""
self.errors_total.labels(error_type=error_type).inc()
def update_positions(self, count: int):
"""Update position count"""
self.positions_count.set(count)
def update_balance(self, balance: float):
"""Update account balance"""
self.account_balance.set(balance)
Telegram Alerts
# alerts.py
import requests
import logging
class AlertManager:
"""
Send alerts via Telegram
"""
def __init__(self, bot_token: str, chat_id: str):
self.bot_token = bot_token
self.chat_id = chat_id
self.logger = logging.getLogger(__name__)
def send_alert(self, message: str, level: str = 'INFO'):
"""Send alert message"""
emoji = {
'INFO': 'ℹ️',
'WARNING': '⚠️',
'ERROR': '🚨',
'SUCCESS': '✅'
}
formatted_message = f"{emoji.get(level, '')} {level}\n\n{message}"
try:
url = f"https://api.telegram.org/bot{self.bot_token}/sendMessage"
payload = {
'chat_id': self.chat_id,
'text': formatted_message,
'parse_mode': 'HTML'
}
response = requests.post(url, json=payload)
response.raise_for_status()
except Exception as e:
self.logger.error(f"Failed to send alert: {e}")
def send_trade_alert(self, trade: Dict):
"""Send trade execution alert"""
message = f"""
<b>Trade Executed</b>
Symbol: {trade['symbol']}
Action: {trade['action']}
Quantity: {trade['quantity']}
Price: ₹{trade['price']}
Order ID: {trade['order_id']}
"""
self.send_alert(message, 'SUCCESS')
def send_error_alert(self, error: str):
"""Send error alert"""
message = f"""
<b>Trading Bot Error</b>
{error}
Check logs for details.
"""
self.send_alert(message, 'ERROR')
Conclusion
Transforming Jupyter notebooks into production systems requires:
- Code refactoring: Functions → Classes → Modules
- Configuration management: Hardcoded values → YAML configs
- Data pipeline: Manual imports → Automated fetching & caching
- Execution engine: Test orders → Production-ready trading
- Infrastructure: Ad-hoc scripts → Systemd services
- Monitoring: No visibility → Comprehensive metrics & alerts
Ready to deploy your strategies? Contact us for professional production deployment services.