AWS Lambda调用Kite WebSocket时触发twisted.internet.error.ReactorNotRestartable错误
问题:AWS Lambda中KiteTicker第二次调用抛出ReactorNotRestartable错误
需求
- 每10分钟运行一次AWS Lambda函数,订阅单只或多只股票的tick数据
- 仅当tick的
last_trade_time秒数处于1-20区间时,执行数据收集逻辑 - 运行指定时长(ttl_seconds)后关闭WebSocket连接并退出函数
原有实现代码
from kiteconnect import KiteTicker from datetime import datetime, timezone import json # 假设这些变量在Lambda环境中已配置 api_key = "YOUR_API_KEY" access_token = "YOUR_ACCESS_TOKEN" i_token_list = [INSTRUMENT_TOKEN_1, INSTRUMENT_TOKEN_2] tick_data = [] unique_tk = [] start_time = datetime.now(timezone("Asia/Kolkata")) instrument_token = i_token_list print("Tick Stock Instrument Token: ", instrument_token) def on_ticks(ws, ticks): tk = {} tk['t_stamp'] = ticks[0].get("last_trade_time") tk['last_price'] = ticks[0].get("last_price") # print("Ticks: ",tk['t_stamp'],tk['last_price']) return None def on_connect(ws, response): ws.subscribe(instrument_token) ws.set_mode(ws.MODE_FULL, instrument_token) def on_close(ws, code, reason): ws.stop() # Assign the callbacks. kws = KiteTicker(api_key, access_token) kws.on_ticks = on_ticks kws.on_connect = on_connect kws.on_close = on_close kws.connect(threaded=True) # kws.connect() while True: def on_ticks(ws, ticks): feed_data(ticks) def feed_data(ticks): for tick in ticks: # print(tick) tk = {} tk['instrument_token'] = tick['instrument_token'] tk['t_stamp'] = tick["last_trade_time"].astimezone(timezone('Asia/Kolkata')).strftime("%Y-%m-%d %H:%M:%S") tk['last_price'] = tick["last_price"] if 0 < tick["last_trade_time"].second < 20: tick_data.append(tk) kws.on_ticks=on_ticks def update_item(l_tick) : print(json.dumps(l_ticks, indent = 2)) if (datetime.now(timezone("Asia/Kolkata")) -start_time).total_seconds() % 2 == 0: # 假设update_dynamo_item是已实现的写入DynamoDB函数 update_dynamo_item(tick_data) ttl_seconds = 20 if (datetime.now(timezone("Asia/Kolkata")) -start_time).total_seconds() > ttl_seconds: update_dynamo_item(tick_data) print(str(ttl_seconds) + " seconds passed, closing connection...") kws.unsubscribe(instrument_token) kws.close() kws.stop() print("Tick Data: ", json.dumps(tick_data, indent=2)) break
错误日志
2025-06-26T09:19:50.536+05:30 Exception in thread Thread-2 (run): 2025-06-26T09:19:51.356+05:30 Traceback (most recent call last): 2025-06-26T09:19:51.416+05:30 File "/var/lang/lib/python3.13/threading.py", line 1041, in _bootstrap_inner 2025-06-26T09:19:51.416+05:30 self.run() 2025-06-26T09:19:51.416+05:30 ~~~~~~~~^^ 2025-06-26T09:19:51.456+05:30 File "/var/lang/lib/python3.13/threading.py", line 992, in run 2025-06-26T09:19:51.456+05:30 self._target(*self._args, **self._kwargs) 2025-06-26T09:19:51.456+05:30 ~~~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ 2025-06-26T09:19:51.476+05:30 File "/opt/python/twisted/internet/base.py", line 695, in run 2025-06-26T09:19:51.476+05:30 self.startRunning(installSignalHandlers=installSignalHandlers) 2025-06-26T09:19:51.476+05:30 ~~~~~~~~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ 2025-06-26T09:19:51.516+05:30 File "/opt/python/twisted/internet/base.py", line 926, in startRunning 2025-06-26T09:19:51.516+05:30 raise error.ReactorNotRestartable() 2025-06-26T09:19:51.536+05:30 twisted.internet.error.ReactorNotRestartable
解决方案
核心原因
KiteTicker依赖Twisted框架的Reactor处理WebSocket连接,而AWS Lambda的执行环境是进程复用的:第一次调用启动的Reactor会留在进程内存中,第二次调用尝试重新启动Reactor时,就会触发ReactorNotRestartable错误。此外原有代码存在循环内重复定义回调函数的问题,进一步加剧了逻辑混乱。
修复后的代码
from kiteconnect import KiteTicker from datetime import datetime, timezone import json import os def lambda_handler(event, context): # 从Lambda环境变量获取配置 api_key = os.environ.get("KITE_API_KEY") access_token = os.environ.get("KITE_ACCESS_TOKEN") i_token_list = list(map(int, os.environ.get("INSTRUMENT_TOKENS").split(","))) ttl_seconds = int(os.environ.get("TTL_SECONDS", 20)) tick_data = [] start_time = datetime.now(timezone("Asia/Kolkata")) def feed_data(ticks): for tick in ticks: tk = { "instrument_token": tick["instrument_token"], "t_stamp": tick["last_trade_time"].astimezone(timezone('Asia/Kolkata')).strftime("%Y-%m-%d %H:%M:%S"), "last_price": tick["last_price"] } # 只收集秒数1-20的tick数据 if 0 < tick["last_trade_time"].second < 20: tick_data.append(tk) def on_ticks(ws, ticks): feed_data(ticks) def on_connect(ws, response): ws.subscribe(i_token_list) ws.set_mode(ws.MODE_FULL, i_token_list) def on_close(ws, code, reason): ws.stop() # 每次调用都创建新的KiteTicker实例,避免复用旧的Reactor kws = KiteTicker(api_key, access_token) kws.on_ticks = on_ticks kws.on_connect = on_connect kws.on_close = on_close # 启动WebSocket连接(线程模式) kws.connect(threaded=True) # 循环等待直到达到超时时间 while True: elapsed = (datetime.now(timezone("Asia/Kolkata")) - start_time).total_seconds() # 每2秒写入一次DynamoDB(可选,按需调整) if elapsed % 2 == 0: update_dynamo_item(tick_data) # 达到超时时间后清理资源 if elapsed > ttl_seconds: update_dynamo_item(tick_data) print(f"{ttl_seconds} seconds passed, closing connection...") kws.unsubscribe(i_token_list) kws.close() kws.stop() print("Tick Data: ", json.dumps(tick_data, indent=2)) break return { "statusCode": 200, "body": json.dumps({"message": "Data collection completed", "tick_count": len(tick_data)}) } # 假设的DynamoDB写入函数,需根据实际实现调整 def update_dynamo_item(data): if not data: return # 此处替换为实际的DynamoDB写入逻辑 print(f"Writing {len(data)} tick records to DynamoDB")
关键修复点
- 将KiteTicker实例移到Lambda处理函数内部:每次Lambda调用都会创建全新的KiteTicker实例和对应的Twisted Reactor,避免复用进程中残留的旧Reactor。
- 移除循环内的重复函数定义:把
on_ticks、feed_data等函数定义放在处理函数顶层,只绑定一次回调,避免逻辑混乱。 - 使用环境变量配置参数:将API密钥、合约代码等敏感信息或可变参数放在Lambda环境变量中,提升代码安全性和可维护性。
- 简化循环逻辑:通过计算已运行时间来控制循环,避免不必要的重复操作。
内容的提问来源于stack exchange,提问作者karthik
相关产品推荐
相关产品推荐

