| 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...") |
| |
| |
| 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...") |
| |
| |
| df = pd.DataFrame(raw_instruments) |
| if df.empty: |
| raise ValueError("Downloaded instruments payload is empty.") |
| |
| |
| df['exchange'] = df['exchange'].astype(str) |
| df['name'] = df['name'].astype(str) |
| |
| |
| 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.") |
| |
| |
| today = datetime.now() |
| year_str = today.strftime("%Y") |
| month_str = today.strftime("%m") |
| day_str = today.strftime("%d") |
| |
| |
| |
| 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}" |
| ) |
| |
| |
| 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}...") |
| |
| |
| 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 |
|
|