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

解决COM4端口竞争问题,实现电力与电流数据的稳定SQL采集

解决COM4端口竞争问题的方案

问题背景

现有Plotly-Dash项目从两个SQLite数据库读取数据,这两个数据库分别由电力数据采集程序和电流数据采集程序写入。最初两个采集程序均通过Modbus协议连接物理USB端口COM4的设备,后因同时访问设备出现冲突,调整为仅由电力数据采集程序连接COM4,采集完成后将电流数据通过Socket发送给电流采集程序处理。但目前仍存在端口竞争问题:触发采集时,进程间抢占COM4端口导致数据返回0或无效值。

终端错误日志

Collecting Data [0]
Collecting Data [0]
3.45.3
3.45.3
Opened database successfully
Table created successfully
Starting on Active Power 1
Connected to COM4
Opened database successfully
Table created successfully
Starting on Active Power 1
Could not connect to COM4
Warning: avg_over_2hsec(25) returned None
[0]Active Power 1 Done
Starting on Active Power 2
Could not connect to COM4
Warning: avg_over_2hsec(27) returned None
[0]Active Power 2 Done
Starting on Active Power 3
Could not connect to COM4
Warning: avg_over_2hsec(29) returned None
[0]Active Power 3 Done
Starting on Active Power Total
Could not connect to COM4
Warning: avg_over_2hsec(43) returned None
[0]Active Power Total Done
Starting on Current 1
Could not connect to COM4
[0]Current 1 Done
Starting on Current 2
Could not connect to COM4
[0]Current 2 Done
Starting on Current 3
Could not connect to COM4
[0]Current 3 Done
Starting on Current Total
Could not connect to COM4
[0]Current Total Done
[0]Sent current data to Process B.

相关代码

电力数据采集程序核心代码

def data_gathering(n):
    with sqlite3.connect('power.db', timeout = 5, check_same_thread = False) as conn:
        conn.execute('PRAGMA journal_mode=WAL;')
        print("Opened database successfully")

        #conn.execute('DROP TABLE IF EXISTS POWER;')

        conn.execute('''
        CREATE TABLE IF NOT EXISTS POWER (
            TIME INT PRIMARY KEY,
            APW1 INT NOT NULL,
            APW2 INT NOT NULL,
            APW3 INT NOT NULL,
            APWTOT INT NOT NULL
        );
        ''')
        conn.commit()

        print("Table created successfully")
        
        time_stamp = int(time.time())
        
        print("Starting on Active Power 1")
        val1 = avg_over_2hsec(25)
        if val1 is None:
            print("Warning: avg_over_2hsec(25) returned None")
            appow1 = 0  
        else:
            appow1 = round(250 * val1)
        
        print(f"[{n}]Active Power 1 Done")

        print("Starting on Active Power 2")
        val2 = avg_over_2hsec(27)
        if val2 is None:
            print("Warning: avg_over_2hsec(27) returned None")
            appow2 = 0 
        else:
            appow2 = round(250 * val2)
        print(f"[{n}]Active Power 2 Done")

        print("Starting on Active Power 3")
        val3 = avg_over_2hsec(29)
        if val3 is None:
            print("Warning: avg_over_2hsec(29) returned None")
            appow3 = 0  
        else:
            appow3 = round(250 * val3)
        print(f"[{n}]Active Power 3 Done")

        print("Starting on Active Power Total")
        valt = avg_over_2hsec(43)
        if valt is None:
            print("Warning: avg_over_2hsec(43) returned None")
            appowtot = 0  
        else:
            appowtot = round(250 * valt)
        print(f"[{n}]Active Power Total Done")
        
        conn.execute("INSERT OR REPLACE INTO POWER (TIME, APW1, APW2, APW3, APWTOT) VALUES (?, ?, ?, ?, ?)",
             (time_stamp, appow1, appow2, appow3, appowtot));
        conn.commit()

        #CURRENT GATHERING PORTION 

        time_stamp = int(time.time())
        
        print("Starting on Current 1")
        value = avg_over_short(17)
        if value is None:
            value = 0  # or some default fallback value
        curr1 = round(250 * value)
        print(f"[{n}]Current 1 Done")

        print("Starting on Current 2")
        value2 = avg_over_short(19)
        if value2 is None:
            value2 = 0  # or some default fallback value
        curr2 = round(250 * value2)
        print(f"[{n}]Current 2 Done")

        print("Starting on Current 3")
        value3 = avg_over_short(21)
        if value3 is None:
            value3 = 0  # or some default fallback value
        curr3 = round(250 * value3)        
        print(f"[{n}]Current 3 Done")

        print("Starting on Current Total")
        valuet = avg_over_short(23)
        if valuet is None:
            valuet = 0  # or some default fallback value
        currtot = round(250 * valuet)
        print(f"[{n}]Current Total Done")

        current_data = {
        'TIME': time_stamp,
        'CURRY1': curr1,
        'CURRY2': curr2,
        'CURRY3': curr3,
        'CURRYTOT': currtot
    }
        send_curr_data(current_data)
        print(f"[{n}]Sent current data to Process B.")


def send_curr_data(current_data):
    try:
        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
            s.connect(('localhost', 9876))  # Connect to process B
            message = json.dumps(current_data)
            s.sendall(message.encode('utf-8'))
    except ConnectionRefusedError:
        print("gatheringcurrentdata.py not running or connection refused")

    except Exception as e:
        print(f"Error: {e}")


def save_to_db(data):
    with sqlite3.connect('current_full.db') as conn:
        cursor = conn.cursor()
        # Insert current_data, with timestamp, CURRY1, CURRY2, 
        cursor.execute("INSERT INTO CURRENT (TIME, CURRY1, CURRY2, CURRY3, CURRYTOT) VALUES (?, ?, ?, ?, ?)",
                       (data['TIME'], data['CURRY1'], data['CURRY2'], data['CURRY3'],     data['CURRYTOT']))
        conn.commit()   

if __name__ == '__main__':
    n = 0
    freeze_support()

    now = datetime.now()
    sec_after_hour = now.minute * 60 + now.second
    sec_til_cleanstart = 600 - (sec_after_hour % 600)
    print("Waiting to start on xy:00,", sec_til_cleanstart, "seconds left")
    time.sleep(sec_til_cleanstart)


    while True:
        nowDT = datetime.now()
        if nowDT.minute % 10 >= 8:
            wait_secs = ((10- nowDT.minute % 10) % 10) * 60 - nowDT.second
            print(f"[DataGathering] Outside write window. Sleeping {wait_secs} seconds.")
            time.sleep(wait_secs)
            continue

        now = time.time()


        print("Collecting Data", [n], flush=True)
        p = Process(target=data_gathering, args=(n,))
        p.start()

        p.join(timeout=480)

        if p.is_alive():
            print("Function timed out and was stopped.", flush=True)
            p.terminate()
            print("Terminate called")
            p.join()
            print("Process joined after termination", flush=True)  
        else:
            print("Data collected in time", flush=True) 
            try:
                with sqlite3.connect("power.db", timeout = 5, check_same_thread = False) as conn:
                    conn.execute("ATTACH DATABASE 'power_full.db' AS power_full")

                    conn.execute('''
                        CREATE TABLE IF NOT EXISTS power_full.POWER (
                            TIME INT PRIMARY KEY,
                            APW1 INT NOT NULL,
                            APW2 INT NOT NULL,
                            APW3 INT NOT NULL,
                            APWTOT INT NOT NULL
                        );
                        ''')
                
                    conn.execute('''
                        INSERT OR IGNORE INTO power_full.POWER
                                SELECT * FROM main.POWER
                    ''')

                    conn.commit()
                    print("Data copied across files")
        
            except Exception as e:#
                print("Failed to copy data across files")
    
            elapsed_time = time.time() - now
            remaining_time = 600 - elapsed_time
            if remaining_time > 0:
                print (f"Sleeping for", [remaining_time])
                time.sleep(remaining_time)
        
        n += 1  

电流数据采集程序核心代码

def run_current_data_server(host='localhost', port=9876):
    server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server.bind((host, port))
    server.listen()
    print(f"[Server] Listening on {host}:{port} for current data...")

    while True:
        conn, addr = server.accept()
        with conn:
            data_bytes = b''
            while True:
                chunk = conn.recv(1024)
                if not chunk:
                    break
                data_bytes += chunk
            try:
                data = json.loads(data_bytes.decode('utf-8'))
                save_to_db(data)
                print(f"[Server] Received and saved current data: {data}")
            except Exception as e:
                print(f"[Server] Error processing data: {e}")

def save_to_db(data):
    with sqlite3.connect('current_full.db') as conn:
        cursor = conn.cursor()
        cursor.execute("""
            CREATE TABLE IF NOT EXISTS CURRENT (
                TIME INT PRIMARY KEY,
                CURRY1 INT NOT NULL,
                CURRY2 INT NOT NULL,
                CURRY3 INT NOT NULL,
                CURRYTOT INT NOT NULL
            );
        """)
        cursor.execute("INSERT OR IGNORE INTO CURRENT (TIME, CURRY1, CURRY2, CURRY3, CURRYTOT)     VALUES (?, ?, ?, ?, ?)",
                       (data['TIME'], data['CURRY1'], data['CURRY2'], data['CURRY3'],     data['CURRYTOT']))
        conn.commit()

if __name__ == '__main__':
    run_current_data_server()

解决方案

1. 单进程内复用Modbus连接

日志显示avg_over_2hsec和avg_over_short每次调用都会尝试连接COM4,导致同一进程内多次抢占端口。修改逻辑,在采集开始时建立一次连接,全程复用:

  • 在data_gathering函数开头创建Modbus连接对象
  • 将连接对象作为参数传递给所有需要访问设备的采集函数
  • 采集完成后统一关闭连接

示例修改:

def data_gathering(n):
    modbus_client = None
    try:
        # 根据实际设备参数初始化连接
        modbus_client = ModbusClient(port='COM4', baudrate=9600, timeout=1)
        modbus_client.connect()
        print("Connected to COM4 successfully")
    except Exception as e:
        print(f"Failed to connect to COM4: {e}")

    with sqlite3.connect('power.db', timeout = 5, check_same_thread = False) as conn:
        # ... 原有数据库操作代码 ...
        
        # 传递连接对象给采集函数
        val1 = avg_over_2hsec(25, modbus_client)
        # ... 其他电力数据采集 ...
        
        value = avg_over_short(17, modbus_client)
        # ... 其他电流数据采集 ...

    # 采集完成后关闭连接
    if modbus_client and modbus_client.is_open():
        modbus_client.close()

同时修改采集函数:

def avg_over_2hsec(register_addr, modbus_client):
    if not modbus_client or not modbus_client.is_open():
        return None
    # 原有采集逻辑,直接使用传入的连接
    # ...

def avg_over_short(register_addr, modbus_client):
    if not modbus_client or not modbus_client.is_open():
        return None
    # 原有采集逻辑,直接使用传入的连接
    # ...

2. 跨进程锁控制端口访问

若需保留多进程结构,使用跨进程锁限制同一时间只有一个进程访问COM4:

  • Windows环境可使用win32file实现文件锁
  • Linux环境可使用fcntl实现文件锁
  • 在启动采集进程前获取锁,采集完成后释放锁

示例修改(Windows):

import win32file
import win32con

def acquire_com_lock():
    lock_path = "com4_lock.lock"
    handle = win32file.CreateFile(
        lock_path,
        win32con.GENERIC_READ | win32con.GENERIC_WRITE,
        0,
        None,
        win32con.CREATE_ALWAYS,
        win32con.FILE_ATTRIBUTE_NORMAL,
        None
    )
    win32file.LockFileEx(handle, win32con.LOCKFILE_EXCLUSIVE_LOCK, 0, 0xffff0000, None)
    return handle

def release_com_lock(handle):
    win32file.UnlockFileEx(handle, 0, 0xffff0000, None)
    win32file.CloseHandle(handle)

# 主循环中使用锁
while True:
    # ... 原有等待逻辑 ...
    
    lock_handle = acquire_com_lock()
    
    p = Process(target=data_gathering, args=(n,))
    p.start()
    p.join(timeout=480)
    
    release_com_lock(lock_handle)
    
    # ... 原有后续逻辑 ...

3. 重构为单进程服务模式

彻底避免多进程竞争,将设备访问逻辑独立为单进程服务:

  • 编写独立的Modbus采集服务,负责从COM4获取所有电力、电流数据
  • 电力和电流采集程序通过Socket或共享内存从该服务获取数据,不再直接访问COM4

此方案从根源解决端口竞争问题,架构更稳定易维护。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 19:17:01