You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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")

关键修复点

  1. 将KiteTicker实例移到Lambda处理函数内部:每次Lambda调用都会创建全新的KiteTicker实例和对应的Twisted Reactor,避免复用进程中残留的旧Reactor。
  2. 移除循环内的重复函数定义:把on_ticks、feed_data等函数定义放在处理函数顶层,只绑定一次回调,避免逻辑混乱。
  3. 使用环境变量配置参数:将API密钥、合约代码等敏感信息或可变参数放在Lambda环境变量中,提升代码安全性和可维护性。
  4. 简化循环逻辑:通过计算已运行时间来控制循环,避免不必要的重复操作。

内容的提问来源于stack exchange,提问作者karthik

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 19:35:04