У меня есть следующий код, который выполняется как подпроцесс.Popen() из основного процесса fastapi.
import asyncio
import datetime
import signal
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from src.app.screens.iv import calculate_iv
from src.app.screens.queries.daily_tasks import (
ALL_OPTIONS_BHAVCOPY,
ALL_OPTIONS_LIVE,
LAST_BHAVCOPY,
LAST_CH_INSERT,
)
from src.config import config
from src.db.databases import jf_conn
from src.app.utils.log import logger
from src.config.constants import (
MARKET_END,
MARKET_START
)
from src.app.utils.holidays import Holiday
from src.db.clickhouse_connect import clickhouse_conn
from src.app.services.update_redis import update_redis
from src.app.api.historical import save_iv_historical_daily
scheduler = AsyncIOScheduler()
stop_event = asyncio.Event()
def handle_shutdown_signal(s, f):
logger.info("Shutting down daily worker.")
stop_event.set()
signal.signal(signal.SIGTERM, lambda s, f: handle_shutdown_signal(s, f))
signal.signal(signal.SIGINT, lambda s, f: handle_shutdown_signal(s, f))
async def daily_task(running_tasks:dict):
try:
if running_tasks.get("daily_task"):
logger.info("Task 'daily_task' already exists. Returning from the scheduled job.")
return
columns = ["Token", "Symbol", "SymbolWithExpiry", "FuturesPrice", "PrevFuturesPrice", "ExpiryDate",
"StrikePrice", "OptionType", "LotSize", "OptionOpen", "OptionHigh", "OptionLow", "Premium", "PrevPremium"]
today = datetime.date.today()
if Holiday.is_holiday():
logger.info('Returning from daily worker. Today is holiday.')
return
while True:
now = datetime.datetime.now()
if MARKET_START None:
straddles = await clickhouse_conn.execute_query(STRADDLE_OPTIONS_QUERY)
straddle_df = pd.DataFrame(straddles, columns=['Symbol', 'ExpiryDate', 'StrikePrice', 'CE', 'CE_change', 'CE_chng_percent',
'PE', 'PE_change', 'PE_chng_percent', 'Straddle', 'StraddleHigh', 'StraddleLow',
'PrevStraddle', 'StraddleChange', 'StraddleChangePercent'])
futures = await clickhouse_conn.execute_query(STRADDLE_FUTURES_QUERY)
futures_df = pd.DataFrame(futures, columns=['Symbol', 'ExpiryDate', 'StrikePrice', 'FuturesPrice', 'LotSize', 'InsDt'])
merged_df = straddle_df.merge(futures_df, on=['Symbol', 'ExpiryDate', 'StrikePrice'])
merged_df['StrikeDiff'] = abs(merged_df['FuturesPrice'] - merged_df['StrikePrice'])
grouped_df = merged_df.groupby(['Symbol', 'ExpiryDate'])
await asyncio.gather(*[individual_straddle_update(symbol, expiry.strftime("%Y-%m-%d"), df)
for (symbol, expiry), df in grouped_df])
async def chart_details_grouped_by_expiry(symbol:str):
data = await clickhouse_conn.execute_query(IV_CHART_DETAILS_SYMBOL, symbol)
if not data:
raise Exception(f"No chart data found for symbol: {symbol}")
df = pd.DataFrame(data, columns=['ExpiryDate', 'iv_ce', 'iv_pe', 'insdt'])
df = df.round(2)
return df.groupby('ExpiryDate')
async def individual_iv_update(symbol, df:pd.DataFrame, ranges) -> None:
results = []
chart_df = await chart_details_grouped_by_expiry(symbol)
for row in df.to_dict(orient='records'):
snapshot = {key: {"CE": {"HighVol": row[f'{key}_IV_CE_HIGH'], "LowVol": row[f'{key}_IV_CE_LOW'], "AvgVol": row[f'{key}_IV_CE_AVG']},
"PE": {"HighVol": row[f'{key}_IV_PE_HIGH'], "LowVol": row[f'{key}_IV_PE_LOW'], "AvgVol": row[f'{key}_IV_PE_AVG']}}
for key in ranges}
call_chart = []
put_chart = []
for r in chart_df.get_group(row['ExpiryDate']).to_dict(orient='records'):
insdt = r['insdt'].strftime("%H:%M:%S")
call_chart.append({'per': r['iv_ce'], 'insDt': insdt})
put_chart.append({'per': r['iv_pe'], 'insDt': insdt})
results.append(({symbol: {
"Symbol": symbol,
"based_on": row['InsDt'].strftime("%Y-%m-%d %H:%M:%S"),
"Volatility": snapshot,
"OptionMetricsCE": {
"StrikePrice": row['StrikePrice'],
"OptionPremium": row['ce_ltp'],
"PriceOfUnderlying": row['Future'],
"DaysToExpiry": row['DaysToExpiry'],
"Type": "CE",
"IV": row['iv_ce'],
"IV Change Percentage": row['iv_ce_changeperc'],
"IV change point": row['iv_ce_change'],
"High": row['iv_ce_high'],
"Low": row['iv_ce_low'],
"Avg": row['iv_ce_avg'],
"Risk-free interest": RISK_FREE_INTEREST_IV,
"Dividend-yield": 0,
},
"OptionMetricsPE": {
"StrikePrice": row['StrikePrice'],
"OptionPremium": row['pe_ltp'],
"PriceOfUnderlying": row['Future'],
"DaysToExpiry": row['DaysToExpiry'],
"Type": "PE",
"IV": row['iv_pe'],
"IV Change Percentage": row['iv_pe_changeperc'],
"IV change point": row['iv_pe_change'],
"High": row['iv_pe_high'],
"Low": row['iv_pe_low'],
"Avg": row['iv_pe_avg'],
"Risk-free interest": RISK_FREE_INTEREST_IV,
"Dividend-yield": 0,
},
"ChartDetails" : {
"CE": call_chart,
"PE": put_chart
}
}}, symbol, row['ExpiryDate'].strftime('%Y-%m-%d')))
return results
async def update_iv() -> None:
today = datetime.date.today()
web_columns = ["Symbol", "ExpiryDate", "StrikePrice", "ce_ltp", "pe_ltp", "Future", "DaysToExpiry",
"iv_ce", "iv_ce_change", "iv_ce_changeperc", "iv_ce_high", "iv_ce_low", "iv_ce_avg",
"iv_pe", "iv_pe_change", "iv_pe_changeperc", "iv_pe_high", "iv_pe_low", "iv_pe_avg", "InsDt"]
web_iv_data = await clickhouse_conn.execute_query(IV_WEB_QUERY)
web_df = pd.DataFrame(web_iv_data, columns=web_columns)
ranges = {
"7Days": (today - datetime.timedelta(days=7), today),
"15Days": (today - datetime.timedelta(days=15), today),
"1Month": (today - datetime.timedelta(days=30), today),
"3Month": (today - datetime.timedelta(days=90), today),
}
cols_list = []
coroutines = []
for label, (start, end) in ranges.items():
coroutines.append(clickhouse_conn.execute_query(IV_SNAPSHOT_QUERY, start, end))
cols_list.append(['Symbol', 'ExpiryDate', f'{label}_IV_CE_HIGH', f'{label}_IV_CE_LOW',
f'{label}_IV_CE_AVG', f'{label}_IV_PE_HIGH', f'{label}_IV_PE_LOW', f'{label}_IV_PE_AVG'])
snapshots = await asyncio.gather(*coroutines)
for i, data in enumerate(snapshots):
df_snap = pd.DataFrame(data, columns=cols_list)
web_df = web_df.merge(df_snap, on=['Symbol', 'ExpiryDate'])
web_df = web_df.round(2)
web_df_grouped = web_df.groupby(['Symbol'])
results = await asyncio.gather(*[individual_iv_update(symbol, df, ranges)
for (symbol,), df in web_df_grouped])
for group in results:
for data, sym, exp in group:
IV_PIPE.set(f"iv_{sym}_{exp}", json.dumps(data, cls=EnhancedJSONEncoder))
await IV_PIPE.execute()
async def update_redis(running_tasks: dict) -> None:
await update_straddle()
await update_iv()
running_tasks.pop('update_redis', None)
Почему значения стрэддла обновляются, а значения iv — нет. Я не получаю журналы ошибок. Каждый день daily_worker запускается в 9:15 утра, и это единственное место, где вызывается эта функция. когда запускается daily_worker, возникает эта проблема, но если я перезапущу приложение, оно будет работать без проблем. daily_worker запускается следующим образом:
в файле main.py, где инициализируется объект приложения fastapi.
@asynccontextmanager
async def lifespan(app: FastAPI):
from src.config import config
from src.db.redis import redis_conn
from src.db.clickhouse_connect import clickhouse_conn
from src.db.databases import sde_conn, jf_conn, js_conn
poetry_env = config.poetry_env
bg_worker = subprocess.Popen([poetry_env, "src/app/services/bg_worker.py"])
scheduler = AsyncIOScheduler()
daily_worker = subprocess.Popen([poetry_env, "src/app/services/daily_worker.py"])
def start_daily_worker():
nonlocal daily_worker
if daily_worker is not None and daily_worker.poll() is None:
daily_worker.terminate()
try:
daily_worker.wait(10)
except subprocess.TimeoutExpired:
daily_worker.kill()
daily_worker = subprocess.Popen([poetry_env, "src/app/services/daily_worker.py"])
scheduler.add_job(start_daily_worker,
'cron',
hour=9,
minute=15,
misfire_grace_time=3600)
scheduler.start()
yield
scheduler.shutdown()
if bg_worker:
bg_worker.terminate()
try:
bg_worker.wait(10)
except subprocess.TimeoutExpired:
bg_worker.kill()
if daily_worker:
daily_worker.terminate()
try:
daily_worker.wait(10)
except subprocess.TimeoutExpired:
daily_worker.kill()
js_conn.disconnect()
await sde_conn.disconnect()
await jf_conn.disconnect()
await clickhouse_conn.disconnect()
await redis_conn.disconnect()