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