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

