首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >元器件车间环境采集:RJ45温湿度变送器TCP/IP断线重连机制代码实现

元器件车间环境采集:RJ45温湿度变送器TCP/IP断线重连机制代码实现

原创
作者头像
盛世宏博小可
发布于 2026-09-24 15:34:11
发布于 2026-09-24 15:34:11
870
举报

元器件车间环境采集:RJ45温湿度变送器TCP/IP断线重连机制代码实现

元器件车间 #RJ45温湿度变送器 #TCP/IP #断线重连 #心跳机制 #指数退避 #socket编程 #嵌入式 #工业以太网 #环境采集

元器件车间的环境有个特点:不是一直平稳运行的。贴片机、回流焊、波峰焊这些大功率设备启停时,电网波动会顺着PoE线路传导,交换机端口可能瞬间出现CRC错误甚至up/down。再加上车间粉尘、振动、温湿度变化,网线接头氧化松动的概率比办公室高得多。所以断线不是"会不会发生"的问题,而是"多久发生一次、发生后多久能恢复"的问题。

先说一个车间现场的真实案例:

某SMT车间,60台RJ45温湿度变送器分布在回流焊区、贴片区、IQC区、仓储区。 回流焊炉启动瞬间,车间东侧12台传感器同时离线,持续8-15秒后自动恢复。 原采集程序用简单while循环,断线后直接close socket再connect,没有退避,没有心跳。 结果:12台同时重连,交换机端口被SYN Flood,恢复时间从8秒拖到2分钟。 优化后:指数退避(1/2/4/8s,上限60s)+ 随机抖动(±30%)+ 心跳保活(30s), 同样工况下,12台在15秒内全部恢复,无交换机拥塞。


一、断线场景分类与应对策略

1. 车间常见断线场景

场景

触发原因

持续时间

发生频率

应对策略

设备启停浪涌

回流焊/贴片机启停,电网波动

数秒至数十秒

每天数次

快速检测 + 指数退避重连

网线接头松动

振动导致RJ45接触不良

间歇性

长期存在

心跳检测 + 自动重连

交换机端口保护

CRC错误累积,端口err-disable

需手动恢复或自动恢复

偶发

检测RST + 告警通知

PoE供电波动

交换机负载突增,PoE预算不足

数秒至数分钟

偶发

检测连接断开 + 退避重连

网络拥塞

车间其他设备大量占带宽

持续

生产高峰期

TCP Keepalive + 超时控制

IP冲突

新设备接入,DHCP分配重复IP

持续

极少

检测RST + 告警

2. 断线检测的三层机制

代码语言:javascript
复制
Layer 1:TCP连接层(最快,毫秒级)
  → send()/recv() 返回错误
  → select()/poll() 检测到socket异常
  → TCP Keepalive 探测对端存活

Layer 2:应用层心跳(秒级)
  → 定期发送心跳请求(如Modbus TCP读取设备状态寄存器)
  → 超时未收到响应判定为断线
  → 心跳间隔:30s(车间建议值)

Layer 3:业务层超时(十秒级)
  → 连续N次采集失败
  → 数据超时未更新
  → 作为最终兜底判断

二、嵌入式端(传感器侧)断线重连实现

1. FreeRTOS + LwIP 架构

代码语言:javascript
复制
/* 传感器固件中的TCP服务器/客户端断线重连任务 */

#include "lwip/tcp.h"
#include "lwip/sockets.h"
#include "freertos/FreeRTOS.h"
#include "freertos/task.h"

#define SERVER_IP       "192.168.10.100"
#define SERVER_PORT     502
#define RECONNECT_BASE_MS  1000    // 基础重连间隔 1s
#define RECONNECT_MAX_MS   60000   // 最大重连间隔 60s
#define KEEPALIVE_IDLE    30       // 30s 无数据发送探测
#define KEEPALIVE_INTVL   5        // 探测间隔 5s
#define KEEPALIVE_PROBES  3        // 连续3次无响应断开
#define HEARTBEAT_INTERVAL 30      // 心跳间隔 30s

static int g_sock = -1;
static bool g_connected = false;

/* 设置TCP Keepalive参数 */
static int set_keepalive(int sockfd)
{
    int optval = 1;
    socklen_t optlen = sizeof(optval);
    
    // 启用TCP Keepalive
    if (setsockopt(sockfd, SOL_SOCKET, SO_KEEPALIVE, &optval, optlen) < 0) {
        LWIP_LOG("setsockopt SO_KEEPALIVE failed");
        return -1;
    }
    
    // 设置Keepalive参数(LwIP需通过TCP层设置)
    // Keepalive空闲时间
    optval = KEEPALIVE_IDLE;
    setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPIDLE, &optval, optlen);
    // Keepalive探测间隔
    optval = KEEPALIVE_INTVL;
    setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPINTVL, &optval, optlen);
    // Keepalive探测次数
    optval = KEEPALIVE_PROBES;
    setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPCNT, &optval, optlen);
    
    return 0;
}

/* 连接服务器 */
static int connect_to_server(void)
{
    struct sockaddr_in server_addr;
    int sock;
    int ret;
    
    sock = socket(AF_INET, SOCK_STREAM, 0);
    if (sock < 0) {
        LWIP_LOG("socket creation failed");
        return -1;
    }
    
    // 设置发送/接收超时
    struct timeval tv;
    tv.tv_sec = 3;
    tv.tv_usec = 0;
    setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
    setsockopt(sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
    
    // 设置Keepalive
    set_keepalive(sock);
    
    memset(&server_addr, 0, sizeof(server_addr));
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(SERVER_PORT);
    server_addr.sin_addr.s_addr = inet_addr(SERVER_IP);
    
    ret = connect(sock, (struct sockaddr *)&server_addr, sizeof(server_addr));
    if (ret < 0) {
        close(sock);
        return -1;
    }
    
    LWIP_LOG("connected to server %s:%d", SERVER_IP, SERVER_PORT);
    g_sock = sock;
    g_connected = true;
    return sock;
}

/* 断开连接 */
static void disconnect_from_server(void)
{
    if (g_sock >= 0) {
        close(g_sock);
        g_sock = -1;
    }
    g_connected = false;
    LWIP_LOG("disconnected from server");
}

/* 指数退避重连(带随机抖动) */
static void reconnect_with_backoff(void)
{
    static uint8_t retry_count = 0;
    uint32_t delay_ms;
    float jitter;
    
    retry_count++;
    
    // 计算指数退避:base * 2^retry,上限 max
    // 1s, 2s, 4s, 8s, 16s, 32s, 60s(cap)...
    if (retry_count >= 6) {
        delay_ms = RECONNECT_MAX_MS;  // 60s cap
    } else {
        delay_ms = RECONNECT_BASE_MS * (1 << (retry_count - 1));
    }
    
    // 添加 ±30% 随机抖动,避免多设备同时重连
    jitter = 0.7f + ((float)rand() / RAND_MAX) * 0.6f;  // 0.7 ~ 1.3
    delay_ms = (uint32_t)(delay_ms * jitter);
    
    LWIP_LOG("reconnect attempt %d, delay %ums (jitter %.2f)", 
             retry_count, delay_ms, jitter);
    
    vTaskDelay(pdMS_TO_TICKS(delay_ms));
    
    disconnect_from_server();
    int sock = connect_to_server();
    if (sock >= 0) {
        retry_count = 0;  // 重连成功,重置计数器
        LWIP_LOG("reconnect success after %d attempts", retry_count);
    }
}

/* TCP通信任务 */
void tcp_comm_task(void *pvParameters)
{
    int sock;
    uint8_t tx_buffer[64];
    uint8_t rx_buffer[64];
    TickType_t last_heartbeat = 0;
    uint32_t heartbeat_count = 0;
    
    // 初始连接
    while (connect_to_server() < 0) {
        LWIP_LOG("initial connection failed, retrying...");
        vTaskDelay(pdMS_TO_TICKS(5000));
    }
    
    while (1) {
        if (!g_connected) {
            // 未连接状态,尝试重连
            reconnect_with_backoff();
            continue;
        }
        
        // 发送心跳(定期发送,保持连接活跃)
        TickType_t now = xTaskGetTickCount();
        if (now - last_heartbeat > pdMS_TO_TICKS(HEARTBEAT_INTERVAL * 1000)) {
            // 构造心跳报文(Modbus TCP读取设备状态)
            build_heartbeat_packet(tx_buffer, heartbeat_count++);
            
            int sent = send(g_sock, tx_buffer, 12, 0);  // MBAP头+功能码
            if (sent < 0) {
                LWIP_LOG("heartbeat send failed, errno=%d", errno);
                g_connected = false;
                continue;
            }
            
            // 等待响应
            int recv_len = recv(g_sock, rx_buffer, sizeof(rx_buffer), 0);
            if (recv_len <= 0) {
                LWIP_LOG("heartbeat recv failed, errno=%d", errno);
                g_connected = false;
                continue;
            }
            
            last_heartbeat = now;
        }
        
        // 正常数据采集和发送
        // ... 采集温湿度数据,封装Modbus TCP报文,send/recv
        
        vTaskDelay(pdMS_TO_TICKS(1000));  // 1s采集周期
    }
}

2. 关键参数说明

参数

值

说明

RECONNECT_BASE_MS

1000ms

基础重连间隔

RECONNECT_MAX_MS

60000ms

最大重连间隔(避免无限增长)

KEEPALIVE_IDLE

30s

TCP连接空闲多久后开始探测

KEEPALIVE_INTVL

5s

探测包发送间隔

KEEPALIVE_PROBES

3

连续3次无响应判定死亡

HEARTBEAT_INTERVAL

30s

应用层心跳间隔

随机抖动

±30%

避免多设备同时重连


三、采集服务器端(Python)断线重连实现

1. 基础版:单设备重连

代码语言:javascript
复制
import socket
import struct
import time
import random
from typing import Optional, Tuple

class SensorTCPClient:
    """RJ45温湿度变送器TCP客户端(带断线重连)"""
    
    def __init__(self, host: str, port: int = 502, 
                 timeout: float = 3.0,
                 reconnect_base: float = 1.0,
                 reconnect_max: float = 60.0,
                 heartbeat_interval: float = 30.0,
                 max_retries: int = 0):  # 0 = 无限重试
        self.host = host
        self.port = port
        self.timeout = timeout
        self.reconnect_base = reconnect_base
        self.reconnect_max = reconnect_max
        self.heartbeat_interval = heartbeat_interval
        self.max_retries = max_retries
        
        self.sock: Optional[socket.socket] = None
        self.connected = False
        self.retry_count = 0
        self.last_heartbeat = 0.0
        self.tid = 0  # Modbus TCP事务ID
        
    def _next_tid(self) -> int:
        self.tid = (self.tid + 1) % 65536
        return self.tid
    
    def _calculate_backoff(self) -> float:
        """计算指数退避延迟(带随机抖动)"""
        if self.retry_count == 0:
            return 0.0
        
        # 指数退避:base * 2^(retry-1)
        exp = min(self.retry_count - 1, 6)  # 最多 2^6 = 64倍
        delay = self.reconnect_base * (2 ** exp)
        delay = min(delay, self.reconnect_max)
        
        # 随机抖动 ±30%
        jitter = 0.7 + random.random() * 0.6  # 0.7 ~ 1.3
        delay *= jitter
        
        return delay
    
    def connect(self) -> bool:
        """建立TCP连接"""
        try:
            self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
            self.sock.settimeout(self.timeout)
            
            # 启用TCP Keepalive
            self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
            
            # Linux下设置Keepalive参数
            if hasattr(socket, 'TCP_KEEPIDLE'):
                self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 30)
            if hasattr(socket, 'TCP_KEEPINTVL'):
                self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 5)
            if hasattr(socket, 'TCP_KEEPCNT'):
                self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 3)
            
            self.sock.connect((self.host, self.port))
            self.connected = True
            self.retry_count = 0
            print(f"[{self.host}] connected")
            return True
            
        except (socket.timeout, socket.error) as e:
            print(f"[{self.host}] connect failed: {e}")
            self.connected = False
            if self.sock:
                self.sock.close()
                self.sock = None
            return False
    
    def disconnect(self):
        """断开连接"""
        self.connected = False
        if self.sock:
            try:
                self.sock.close()
            except:
                pass
            self.sock = None
        print(f"[{self.host}] disconnected")
    
    def reconnect(self) -> bool:
        """断线重连(带指数退避)"""
        self.disconnect()
        
        self.retry_count += 1
        if self.max_retries > 0 and self.retry_count > self.max_retries:
            print(f"[{self.host}] max retries ({self.max_retries}) exceeded")
            return False
        
        delay = self._calculate_backoff()
        if delay > 0:
            print(f"[{self.host}] reconnect #{self.retry_count}, "
                  f"delay {delay:.1f}s")
            time.sleep(delay)
        
        return self.connect()
    
    def send_modbus_request(self, slave_id: int, func_code: int, 
                           reg_addr: int, reg_count: int) -> Optional[bytes]:
        """发送Modbus TCP请求"""
        if not self.connected or not self.sock:
            return None
        
        tid = self._next_tid()
        
        # 构造MBAP头 + PDU
        # MBAP: TID(2) + PID(2) + Length(2) + UnitID(1)
        # PDU: FuncCode(1) + StartAddr(2) + RegCount(2)
        length = 6  # UnitID(1) + PDU(5)
        mbap = struct.pack('>HHHB', tid, 0x0000, length, slave_id)
        pdu = struct.pack('>BHH', func_code, reg_addr, reg_count)
        request = mbap + pdu
        
        try:
            self.sock.sendall(request)
            response = self.sock.recv(256)
            
            if len(response) < 9:
                print(f"[{self.host}] response too short: {len(response)} bytes")
                return None
            
            # 解析MBAP头
            resp_tid, resp_pid, resp_len, resp_uid = struct.unpack('>HHHB', 
                                                                    response[:7])
            
            # 检查异常响应
            func_resp = response[7]
            if func_resp == func_code + 0x80:
                err_code = response[8]
                print(f"[{self.host}] modbus exception: 0x{err_code:02x}")
                return None
            
            return response
            
        except socket.timeout:
            print(f"[{self.host}] send/recv timeout")
            self.connected = False
            return None
        except (ConnectionResetError, BrokenPipeError, OSError) as e:
            print(f"[{self.host}] connection error: {e}")
            self.connected = False
            return None
    
    def read_temperature_humidity(self, slave_id: int = 1,
                                   temp_addr: int = 0, 
                                   hum_addr: int = 1) -> Optional[Tuple[float, float]]:
        """读取温湿度"""
        # 读取2个寄存器(温度+湿度)
        response = self.send_modbus_request(slave_id, 0x03, temp_addr, 2)
        if not response:
            return None
        
        # 解析响应:MBAP(7) + FuncCode(1) + ByteCount(1) + Reg1(2) + Reg2(2)
        byte_count = response[8]
        if byte_count != 4:
            print(f"[{self.host}] unexpected byte count: {byte_count}")
            return None
        
        reg1, reg2 = struct.unpack('>HH', response[9:13])
        temperature = reg1 * 0.1
        humidity = reg2 * 0.1
        
        return temperature, humidity
    
    def heartbeat(self) -> bool:
        """发送心跳(读取设备状态寄存器)"""
        response = self.send_modbus_request(1, 0x03, 0, 1)  # 读寄存器0
        return response is not None
    
    def run(self, interval: float = 5.0):
        """主循环"""
        print(f"[{self.host}] starting client...")
        
        # 初始连接
        while not self.connect():
            delay = self._calculate_backoff()
            print(f"[{self.host}] initial connect failed, retry in {delay:.1f}s")
            if delay > 0:
                time.sleep(delay)
        
        while True:
            try:
                # 心跳检测
                now = time.time()
                if now - self.last_heartbeat > self.heartbeat_interval:
                    if not self.heartbeat():
                        print(f"[{self.host}] heartbeat failed, reconnecting...")
                        if not self.reconnect():
                            continue
                    else:
                        self.last_heartbeat = now
                
                # 读取温湿度
                if self.connected:
                    result = self.read_temperature_humidity()
                    if result:
                        temp, hum = result
                        print(f"[{self.host}] T={temp:.1f}℃ H={hum:.1f}%RH")
                    else:
                        print(f"[{self.host}] read failed, reconnecting...")
                        self.reconnect()
                
                time.sleep(interval)
                
            except KeyboardInterrupt:
                print(f"[{self.host}] stopping...")
                self.disconnect()
                break
            except Exception as e:
                print(f"[{self.host}] unexpected error: {e}")
                self.reconnect()


# 使用示例
if __name__ == "__main__":
    client = SensorTCPClient(
        host="192.168.100.50",
        port=502,
        timeout=3.0,
        reconnect_base=1.0,
        reconnect_max=60.0,
        heartbeat_interval=30.0
    )
    client.run(interval=5.0)

2. 进阶版:多设备并发管理

代码语言:javascript
复制
import asyncio
import socket
import struct
import time
import random
from typing import Dict, Optional, Tuple, List
from dataclasses import dataclass, field
from datetime import datetime
import csv

@dataclass
class SensorConfig:
    """传感器配置"""
    ip: str
    port: int = 502
    slave_id: int = 1
    name: str = ""
    location: str = ""  # 车间区域
    temp_reg: int = 0
    hum_reg: int = 1
    scale: float = 0.1

@dataclass
class SensorState:
    """传感器运行时状态"""
    config: SensorConfig
    connected: bool = False
    retry_count: int = 0
    last_success: float = 0.0
    last_heartbeat: float = 0.0
    consecutive_failures: int = 0
    total_reads: int = 0
    total_failures: int = 0
    last_temp: float = 0.0
    last_hum: float = 0.0
    last_error: str = ""

class MultiSensorManager:
    """多传感器并发管理器"""
    
    def __init__(self, sensors: List[SensorConfig],
                 max_concurrent: int = 50,
                 reconnect_base: float = 1.0,
                 reconnect_max: float = 60.0,
                 heartbeat_interval: float = 30.0,
                 read_timeout: float = 3.0,
                 max_failures: int = 3):
        self.sensors: Dict[str, SensorState] = {
            s.ip: SensorState(config=s) for s in sensors
        }
        self.max_concurrent = max_concurrent
        self.reconnect_base = reconnect_base
        self.reconnect_max = reconnect_max
        self.heartbeat_interval = heartbeat_interval
        self.read_timeout = read_timeout
        self.max_failures = max_failures
        self.semaphore = asyncio.Semaphore(max_concurrent)
        self.tid_counter = 0
        self.running = False
        self.data_callback = None  # 数据回调函数
        
    def _next_tid(self) -> int:
        self.tid_counter = (self.tid_counter + 1) % 65536
        return self.tid_counter
    
    def _calculate_backoff(self, retry_count: int) -> float:
        """指数退避 + 随机抖动"""
        if retry_count == 0:
            return 0.0
        exp = min(retry_count - 1, 6)
        delay = self.reconnect_base * (2 ** exp)
        delay = min(delay, self.reconnect_max)
        jitter = 0.7 + random.random() * 0.6
        return delay * jitter
    
    async def _connect_sensor(self, state: SensorState) -> Optional[asyncio.StreamReader]:
        """建立异步TCP连接"""
        try:
            reader, writer = await asyncio.wait_for(
                asyncio.open_connection(state.config.ip, state.config.port),
                timeout=self.read_timeout
            )
            state.connected = True
            state.retry_count = 0
            state.consecutive_failures = 0
            print(f"[{state.config.ip}] connected")
            return reader, writer
        except (asyncio.TimeoutError, OSError) as e:
            state.connected = False
            state.retry_count += 1
            state.last_error = str(e)
            print(f"[{state.config.ip}] connect failed: {e}")
            return None
    
    async def _read_sensor(self, state: SensorState) -> Optional[Tuple[float, float]]:
        """读取单个传感器"""
        async with self.semaphore:
            if not state.connected:
                return None
            
            try:
                reader, writer = await self._connect_sensor(state)
                if not reader:
                    return None
                
                # 构造Modbus TCP请求
                tid = self._next_tid()
                mbap = struct.pack('>HHHB', tid, 0, 6, state.config.slave_id)
                pdu = struct.pack('>BHH', 0x03, state.config.temp_reg, 2)
                request = mbap + pdu
                
                writer.write(request)
                await writer.drain()
                
                # 等待响应
                response = await asyncio.wait_for(
                    reader.read(256),
                    timeout=self.read_timeout
                )
                
                if len(response) < 9:
                    raise ValueError(f"short response: {len(response)} bytes")
                
                # 解析
                func = response[7]
                if func == 0x83:
                    raise ValueError(f"modbus exception: {response[8]}")
                
                byte_count = response[8]
                if byte_count != 4:
                    raise ValueError(f"unexpected byte count: {byte_count}")
                
                reg1, reg2 = struct.unpack('>HH', response[9:13])
                temp = reg1 * state.config.scale
                hum = reg2 * state.config.scale
                
                state.last_success = time.time()
                state.total_reads += 1
                state.consecutive_failures = 0
                state.last_temp = temp
                state.last_hum = hum
                
                writer.close()
                await writer.wait_closed()
                
                return temp, hum
                
            except asyncio.TimeoutError:
                state.consecutive_failures += 1
                state.total_failures += 1
                state.last_error = "timeout"
                return None
            except (ConnectionResetError, BrokenPipeError, OSError) as e:
                state.consecutive_failures += 1
                state.total_failures += 1
                state.last_error = str(e)
                state.connected = False
                return None
            except Exception as e:
                state.consecutive_failures += 1
                state.total_failures += 1
                state.last_error = str(e)
                return None
    
    async def _monitor_sensor(self, ip: str):
        """单个传感器的监控循环"""
        state = self.sensors[ip]
        
        while self.running:
            try:
                # 检查是否需要重连
                if not state.connected:
                    delay = self._calculate_backoff(state.retry_count)
                    if delay > 0:
                        await asyncio.sleep(delay)
                    # 尝试重连
                    conn = await self._connect_sensor(state)
                    if not conn:
                        continue
                    # 关闭测试连接
                    reader, writer = conn
                    writer.close()
                    await writer.wait_closed()
                
                # 读取数据
                result = await self._read_sensor(state)
                if result:
                    temp, hum = result
                    # 调用数据回调
                    if self.data_callback:
                        await self.data_callback(state, temp, hum)
                else:
                    # 连续失败达到阈值,标记离线
                    if state.consecutive_failures >= self.max_failures:
                        state.connected = False
                        print(f"[{ip}] marked offline after "
                              f"{state.consecutive_failures} failures")
                
                # 心跳检测
                now = time.time()
                if now - state.last_heartbeat > self.heartbeat_interval:
                    if state.connected and state.consecutive_failures == 0:
                        state.last_heartbeat = now
                    elif not state.connected:
                        # 尝试重连
                        delay = self._calculate_backoff(state.retry_count)
                        if delay > 0:
                            await asyncio.sleep(delay)
                        await self._connect_sensor(state)
                
                await asyncio.sleep(5.0)  # 采集间隔
                
            except Exception as e:
                print(f"[{ip}] monitor error: {e}")
                state.last_error = str(e)
                await asyncio.sleep(5.0)
    
    async def start(self, data_callback=None):
        """启动所有传感器监控"""
        self.running = True
        self.data_callback = data_callback
        
        print(f"Starting {len(self.sensors)} sensors...")
        
        # 创建所有监控任务
        tasks = []
        for ip in self.sensors.keys():
            task = asyncio.create_task(self._monitor_sensor(ip))
            tasks.append(task)
            # 错峰启动,避免同时连接
            await asyncio.sleep(random.uniform(0.1, 0.5))
        
        # 等待所有任务
        await asyncio.gather(*tasks, return_exceptions=True)
    
    async def stop(self):
        """停止所有监控"""
        self.running = False
        print("Stopping all sensors...")
    
    def get_status(self) -> Dict:
        """获取所有传感器状态"""
        status = {}
        for ip, state in self.sensors.items():
            status[ip] = {
                'name': state.config.name,
                'location': state.config.location,
                'connected': state.connected,
                'retry_count': state.retry_count,
                'consecutive_failures': state.consecutive_failures,
                'total_reads': state.total_reads,
                'total_failures': state.total_failures,
                'last_temp': state.last_temp,
                'last_hum': state.last_hum,
                'last_error': state.last_error,
                'uptime': time.time() - state.last_success if state.last_success > 0 else 0
            }
        return status

# 使用示例
async def on_data(state: SensorState, temp: float, hum: float):
    """数据回调函数"""
    print(f"[{datetime.now().isoformat()}] {state.config.name} "
          f"({state.config.location}): T={temp:.1f}℃ H={hum:.1f}%RH")

async def main():
    # 从配置文件加载传感器列表
    sensors = [
        SensorConfig(ip="192.168.100.50", name="SMT-01", location="SMT贴片区"),
        SensorConfig(ip="192.168.100.51", name="SMT-02", location="SMT贴片区"),
        SensorConfig(ip="192.168.100.52", name="Reflow-01", location="回流焊区"),
        SensorConfig(ip="192.168.100.53", name="Wave-01", location="波峰焊区"),
        SensorConfig(ip="192.168.100.54", name="IQC-01", location="IQC区"),
        SensorConfig(ip="192.168.100.55", name="Storage-01", location="仓储区"),
    ]
    
    manager = MultiSensorManager(
        sensors=sensors,
        max_concurrent=50,
        reconnect_base=1.0,
        reconnect_max=60.0,
        heartbeat_interval=30.0,
        read_timeout=3.0,
        max_failures=3
    )
    
    try:
        await manager.start(data_callback=on_data)
    except KeyboardInterrupt:
        await manager.stop()
    
    # 输出最终状态
    status = manager.get_status()
    for ip, info in status.items():
        print(f"{ip}: {info}")

if __name__ == "__main__":
    asyncio.run(main())

四、关键设计要点

1. 指数退避 + 随机抖动

代码语言:javascript
复制
为什么需要指数退避?
  → 避免"惊群效应":多台设备同时断线后同时重连,交换机端口拥塞
  → 给网络恢复留出时间:短暂的网络抖动,等待比立即重连更有效

为什么需要随机抖动?
  → 即使退避时间相同,设备启动时间相同也会导致同步
  → 随机抖动让重连时间分散,避免集中冲击

退避序列示例(base=1s, max=60s, jitter=±30%):
  第1次: 0.7s ~ 1.3s
  第2次: 1.4s ~ 2.6s
  第3次: 2.8s ~ 5.2s
  第4次: 5.6s ~ 10.4s
  第5次: 11.2s ~ 20.8s
  第6次: 22.4s ~ 41.6s
  第7次+: 42s ~ 60s (capped)

2. 心跳机制

代码语言:javascript
复制
心跳的目的:
  1. 检测连接是否存活(TCP Keepalive可能不够及时)
  2. 保持连接活跃(防止中间设备/NAT超时断开)
  3. 触发重连(心跳失败 = 连接已死)

心跳间隔选择:
  → 太短:占用带宽和CPU
  → 太长:故障发现延迟大
  → 建议:30s(车间环境)

心跳报文选择:
  → Modbus TCP: 读设备状态寄存器(功能码0x03,寄存器0)
  → 或者自定义心跳命令

3. 连接数限制

代码语言:javascript
复制
传感器端:
  → 最大并发连接数:8(LwIP MEMP_NUM_TCP_PCB)
  → 超过限制的新连接被拒绝(返回RST)

采集端:
  → 并发采集数限制:50-150(Semaphore)
  → 避免采集服务器资源耗尽
  → 避免交换机端口拥塞

五、调试与排障

1. 常见断线原因与排查

现象

可能原因

排查方法

连接立即被RST

传感器IP/端口错误、防火墙拦截

telnet测试、Wireshark抓包

连接超时

网络不可达、传感器离线

ping测试、检查交换机端口

发送成功但无响应

传感器Modbus TCP未启用

Web界面确认、Modbus Poll测试

偶发断开

网络抖动、电磁干扰

检查交换机CRC错误、增加重连退避

批量同时断开

交换机重启、PoE供电波动

检查交换机日志、PoE状态

长时间无法重连

IP冲突、传感器死机

检查IP、重启传感器

2. Wireshark过滤

代码语言:javascript
复制
tcp.flags.reset == 1          # 连接被重置
tcp.analysis.retransmission   # TCP重传
tcp.keepalive                 # Keepalive探测包
ip.addr == 192.168.100.50 and tcp.port == 502  # 指定传感器

六、验收标准

项目

标准

测试方法

断线检测时间

<5s(心跳间隔内)

拔网线,记录检测到的时间

重连成功率

100%(网络恢复后)

断网30s后恢复,检查自动重连

重连时间

<30s(指数退避上限)

记录从断线到恢复的时间

批量重连

无交换机拥塞

同时断开50台,观察恢复情况

数据连续性

断线期间数据标记离线

检查数据库中的离线标记

长期稳定性

7×24h无内存泄漏

持续运行,监控内存和连接数


七、一句话总结

元器件车间RJ45温湿度变送器的断线重连,核心不是"断了能重连",而是"怎么重连才不把网络搞崩"。 指数退避解决"别一起冲"的问题,随机抖动解决"别同步"的问题,心跳机制解决"别等太久才发现"的问题。 嵌入式端和采集端都要实现重连逻辑——传感器要能主动连上来,采集端也要能自动重连上去。 最终交付的不是一段代码,而是一套在网络各种作妖情况下都能自动恢复的通信机制。


关键词:元器件车间、RJ45温湿度变送器、TCP/IP、断线重连、心跳机制、指数退避、随机抖动、socket编程、FreeRTOS、LwIP、Python asyncio、Modbus TCP、Keepalive、多设备并发

标签:#元器件车间 #RJ45温湿度变送器 #TCP/IP #断线重连 #心跳机制 #指数退避 #socket编程 #嵌入式 #FreeRTOS #LwIP #Python #asyncio #ModbusTCP #工业以太网 #环境采集 #动环监控

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 元器件车间环境采集:RJ45温湿度变送器TCP/IP断线重连机制代码实现
  • 元器件车间 #RJ45温湿度变送器 #TCP/IP #断线重连 #心跳机制 #指数退避 #socket编程 #嵌入式 #工业以太网 #环境采集
    • 一、断线场景分类与应对策略
      • 1. 车间常见断线场景
      • 2. 断线检测的三层机制
    • 二、嵌入式端(传感器侧)断线重连实现
      • 1. FreeRTOS + LwIP 架构
      • 2. 关键参数说明
    • 三、采集服务器端(Python)断线重连实现
      • 1. 基础版:单设备重连
      • 2. 进阶版:多设备并发管理
    • 四、关键设计要点
      • 1. 指数退避 + 随机抖动
      • 2. 心跳机制
      • 3. 连接数限制
    • 五、调试与排障
      • 1. 常见断线原因与排查
      • 2. Wireshark过滤
    • 六、验收标准
    • 七、一句话总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档