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
相关产品推荐
相关产品推荐

