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

Threading Loop定时重启函数时,如何确保TCP Socket自动关闭?

解决TCP Socket未关闭导致的资源泄漏问题

看起来你遇到了典型的线程资源泄漏问题——每次threading.Timer启动新的iothub_client_sample_run实例时,旧的线程还在后台运行,它的TCP Socket连接一直处于打开状态,没有被正确关闭,积累多了就会引发各种网络相关的报错。

下面是针对你的问题的完整解决方案,核心思路是让旧线程能够优雅退出并清理Socket资源,同时确保每次新启动的线程都使用全新的Socket连接:

import socket
import select
import time
from threading import Timer
# 假设你的iothub_client、telemetry等相关导入已存在

# 全局标志:通知当前运行的线程停止工作
STOP_CURRENT_RUN = False
# 保存当前Timer实例,避免重复触发
current_timer = None

def iothub_client_sample_run():
    global STOP_CURRENT_RUN, current_timer
    
    # 重置终止标志,当前线程开始正常运行
    STOP_CURRENT_RUN = False
    
    # 取消上一次的Timer(如果存在),防止重复启动
    if current_timer:
        current_timer.cancel()
    
    # 提前安排下一次运行,避免遗漏
    current_timer = Timer(10.0, iothub_client_sample_run)
    current_timer.start()

    try:
        client = iothub_client_init()
        if client.protocol == IoTHubTransportProvider.MQTT:
            print("IoTHubClient is reporting state")
            reported_state = "{\"newState\":\"standBy\"}"
            client.send_reported_state(reported_state, len(reported_state), send_reported_state_callback, SEND_REPORTED_STATE_CONTEXT)
            telemetry.send_telemetry_data(parse_iot_hub_name(), EVENT_SUCCESS, "IoT hub connection is established")
        
        # 将Socket相关逻辑封装为单独函数,便于管理连接生命周期
        def run_socket_client():
            conn = None
            try:
                ip = 'xxxx'
                port = 'xxxx'
                conn = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
                conn.connect((ip, port))
                print('connecting...')
                
                # 用终止标志替代死循环,让线程能响应停止信号
                while not STOP_CURRENT_RUN:
                    try:
                        ready_to_read, ready_to_write, in_error = select.select([conn,], [conn,], [], 5)
                    except:
                        print('connection error')
                        break
                    
                    if len(ready_to_read) > 0:
                        recv2 = conn.recv(2048)
                        if not recv2:  # 对方关闭连接时主动退出
                            break
                        
                        try:
                            recv3 = recv2.split(',')
                            date = float(recv3[2])
                            time_val = float(recv3[3])  # 避免与time模块重名
                            heave= float(recv3[4])
                            north= float(recv3[5])
                            east = float(recv3[6])
                            msg_txt_formatted = MSG_TXT % (date, time_val, heave, north, east)
                            
                            print(msg_txt_formatted)
                            message = IoTHubMessage(msg_txt_formatted)
                            message.message_id = "message_%d" % MESSAGE_COUNT
                            message.correlation_id = "correlation_%d" % MESSAGE_COUNT
                            prop_map = message.properties()
                            
                            client.send_event_async(message, send_confirmation_callback, MESSAGE_COUNT)
                            print("IoTHubClient.send_event_async accepted message [%d] for transmission to IoT Hub." % MESSAGE_COUNT)
                            status = client.get_send_status()
                            print("Send status: %s" % status)
                            
                            global MESSAGE_COUNT, MESSAGE_SWITCH
                            MESSAGE_COUNT += 1
                            time.sleep(config.MESSAGE_TIMESPAN / 4000.0)
                            
                        except IoTHubError as iothub_error:
                            print("Unexpected error %s from IoTHub" % iothub_error)
                            telemetry.send_telemetry_data(parse_iot_hub_name(), EVENT_FAILED, "Unexpected error %s from IoTHub" % iothub_error)
                            return
                        except Exception as e:
                            print(f"Data processing error: {e}")
            except Exception as e:
                print(f'Re-connecting due to error: {e}')
            finally:
                # 无论正常退出还是异常,都确保Socket被关闭
                if conn:
                    try:
                        conn.close()
                        print("Socket connection closed gracefully")
                    except:
                        pass
        
        # 运行Socket客户端,直到收到停止信号或MESSAGE_SWITCH关闭
        while not STOP_CURRENT_RUN and MESSAGE_SWITCH:
            run_socket_client()
            # 连接断开后短暂等待再重连,避免频繁重试
            if not STOP_CURRENT_RUN:
                time.sleep(1)
                
    except KeyboardInterrupt:
        print("IoTHubClient sample stopped")
        # 终止所有运行中的线程并清理Timer
        STOP_CURRENT_RUN = True
        if current_timer:
            current_timer.cancel()

# 启动第一个运行实例
iothub_client_sample_run()

关键修改说明

  • STOP_CURRENT_RUN终止标志:用来通知旧线程退出所有循环,触发finally块关闭Socket,避免资源泄漏。
  • Socket逻辑封装+try...finally:将Socket的创建、连接、接收逻辑独立封装,确保无论正常退出还是发生异常,Socket都会被关闭。
  • 替换while True为while not STOP_CURRENT_RUN:让线程能够响应终止信号,不再无限挂起。
  • Timer实例管理:保存当前Timer的引用,每次启动新Timer前取消旧的,防止重复触发多个线程。
  • Socket断开处理:增加if not recv2判断,当对方主动关闭连接时,主动退出循环并关闭Socket。

额外注意事项

  • 如果需要更严谨的线程管理,可以用类封装全局变量,避免全局变量的副作用。
  • 若有其他线程修改MESSAGE_SWITCH,建议用threading.Lock保证线程安全。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:40:36