BQE / engine /instrument_parser.py
DevWizard-Vandan
Deploy BharatQuant Engine (BQE) to Hugging Face Spaces
2418470
Raw
History Blame Contribute Delete
3.65 kB
import os
import pathlib
import logging
import asyncio
import pandas as pd
from datetime import datetime
from typing import Optional
from core.storage_manager import StorageManager
from connectors.kite_client import KiteBrokerClient
logger = logging.getLogger("bqe.engine.instrument_parser")
class MarketInstrumentParser:
"""
Downloads daily master option instrument lists from Zerodha Kite,
filters them for liquid index options (NIFTY & BANKNIFTY), and archives
them as Parquet files partitioned by date in Hive structure.
"""
def __init__(self, storage_manager: Optional[StorageManager] = None):
"""
Initializes the parser with BQE storage structures and broker clients.
"""
self.storage_manager = storage_manager or StorageManager()
self.broker_client = KiteBrokerClient()
async def download_and_archive_instruments(self) -> pathlib.Path:
"""
Retrieves all instruments, filters for NFO liquid indices,
and saves the snapshot to a partitioned Parquet file.
Returns:
pathlib.Path: Absolute path to the saved Parquet file.
"""
if not self.broker_client.kite_connect:
raise ValueError("No active Kite session. Authentication is required.")
logger.info("Fetching master instruments from Zerodha...")
# Offload blocking network call to separate thread
loop = asyncio.get_event_loop()
raw_instruments = await loop.run_in_executor(
None,
lambda: self.broker_client.kite_connect.instruments()
)
logger.info(f"Downloaded {len(raw_instruments)} total instruments. Processing...")
# Convert to Pandas DataFrame
df = pd.DataFrame(raw_instruments)
if df.empty:
raise ValueError("Downloaded instruments payload is empty.")
# Standardize exchange and name column types
df['exchange'] = df['exchange'].astype(str)
df['name'] = df['name'].astype(str)
# Filter for NFO index options (NIFTY & BANKNIFTY)
filtered_df = df[
(df['exchange'] == 'NFO') &
(df['name'].isin(['NIFTY', 'BANKNIFTY']))
].copy()
logger.info(f"Filtered down to {len(filtered_df)} liquid NFO derivative instruments.")
# Determine today's date in 2026
today = datetime.now()
year_str = today.strftime("%Y")
month_str = today.strftime("%m")
day_str = today.strftime("%d")
# Construct Hive partition path layout:
# data_lake/derivatives/instruments/year=YYYY/month=MM/day=DD/master_instruments.parquet
partition_dir = os.path.join(
self.storage_manager.data_lake_dir,
"derivatives",
"instruments",
f"year={year_str}",
f"month={month_str}",
f"day={day_str}"
)
# Ensure directories exist
pathlib.Path(partition_dir).mkdir(parents=True, exist_ok=True)
target_path = pathlib.Path(os.path.join(partition_dir, "master_instruments.parquet")).resolve()
logger.info(f"Writing derivative master snapshot to: {target_path}...")
# Offload blocking write to thread pool
await loop.run_in_executor(
None,
lambda: filtered_df.to_parquet(str(target_path), engine='pyarrow', index=False)
)
logger.info("Derivative master instruments snapshot archived successfully.")
return target_path