向Queue添加List时数据丢失问题排查求助
我搭建了一套通过Queue实现数据库写入单点操作的系统:从queue中取出记录添加到list,当list中的记录达到指定数量后,将其推送至writer_queue,由另一个thread负责写入SQLite数据库。但发现将list添加到writer_queue时,部分记录丢失,导致最终数据库表存在数据缺口。
问题似乎出在Tables.insert_batch_data()与Tables.add_to_writer_queue()之间,传递到add_to_writer_queue()的List实际长度常与传入的batch长度不符。我未在文档中找到Queue传递数据的总量限制,困惑数据丢失的原因及如何确保数据完整传输。
以下是相关代码:
import os import time import sqlite3 from queue import Queue from pydantic import BaseModel from typing import List, Dict, Optional from threading import Thread from dataclasses import dataclass class Tables(BaseModel): max_buffer_length: int = 1000 rt_table_name: str = '' query_schemas: Dict = { 'RawTable': ['Time', 'Position', 'RPM', 'Flow', 'Density', 'Pressure', 'Tension', 'Torque', 'Weight',], } table_schemas: Dict = { 'RawTable': ['time', 'position', 'rpm', 'flow', 'density', 'pressure', 'tension', 'torque', 'weight',], } def insert_data(self, buffer: Optional[Queue] = None, db_path: Optional[str] = None, writer_queue: Optional[Queue] = None): conn = sqlite3.connect(db_path) table_name = self.rt_table_name cursor = conn.cursor() try: if buffer: self.insert_batch_data(conn, cursor, buffer, table_name, writer_queue) except sqlite3.Error as e: print(f"An error occurred: {e}") conn.rollback() finally: conn.commit() cursor.close() conn.close() def insert_batch_data(self, buffer: Queue, table_name: str, writer_queue: Queue): # Insert data query = self.get_query(table_name) batch = [] while True: if buffer.empty(): time.sleep(5) continue item = buffer.get() if item is None: print("Reached Sentinal Value... Exiting thread...") break batch.append(item) if len(batch) == self.max_buffer_length: self.add_to_writer_queue(query, table_name, batch, writer_queue) print(f"Number of items added to writer_queue: {len(batch)}") batch.clear() # Insert any remaining records in the batch if batch: self.add_to_writer_queue(query, table_name, batch, writer_queue) def get_query(self, table_name: str) -> str: if not table_name: raise ValueError("Table name must not be empty") columns = self.query_schemas[table_name] placeholders = ', '.join(['?' for _ in columns]) query = f"INSERT INTO {table_name} ({', '.join(columns)}) VALUES ({placeholders})" return query def insert_records(self, cursor: sqlite3.Cursor, conn: sqlite3.Connection, query: str, batch: List, table_name: str): try: columns = self.table_schemas[table_name] data_tuples = [ tuple(getattr(row, col) for col in columns) for row in batch ] cursor.executemany(query, data_tuples) conn.commit() print(f"Inserted {len(batch)} records into {table_name}") except sqlite3.Error as e: print(f"SQLite error occurred while inserting records into {table_name}: {e}") conn.rollback() except Exception as e: print(f"Unexpected error occurred while inserting records into {table_name}: {e}") conn.rollback() def process_db_writes(self, writer_queue: Queue, db_path: str): try: conn = sqlite3.connect(db_path) cursor = conn.cursor() while True: while writer_queue.empty(): time.sleep(2) # Sleep for 30 seconds query, table_name, data = writer_queue.get() assert len(data) == self.max_buffer_length, f"Expected {self.max_buffer_length} items, received: {len(data)}\n" self.insert_records(cursor, conn, query, data, table_name) except Exception as e: print(f"Error encountered in process_db_writes: {str(e)}") finally: cursor.close() conn.close() def add_to_writer_queue(self, query: str, table_name: str, batch: List, writer_queue: Queue): while writer_queue.full(): time.sleep(1) assert len(batch) == self.max_buffer_length, f"Expected {self.max_buffer_length} items, received: {len(batch)}\n" writer_queue.put((query, table_name, batch)) @dataclass class RawData: time: float position: float = 0.0 rpm: float = 0.0 flow: float = 0.0 density: float = 0.0 pressure: float = 0.0 tension: float = 0.0 torque: float = 0.0 weight: float = 0.0 class Raw(Tables): def __init__(self, **data): super().__init__(**data) self.rt_table_name = 'RawTable' def populate_queue(buffer: Queue): for i in range(1_000_000): while buffer.full(): time.sleep(1) buffer.put(RawData(time=i)) def fetch_raw_data(db_path: str, rt_db_path: str, db_writer: Queue, max_buffer_length: int, ): try: conn = sqlite3.connect(db_path) buffer = Queue(maxsize=max_buffer_length) raw = Raw(max_buffer_length=max_buffer_length) # Starting other threads fetcher = Thread(target=raw.populate_queue, args=(buffer)) writer = Thread(target=raw.insert_data, args=(buffer, rt_db_path, db_writer)) # Start the Threads fetcher.start() writer.start() # Join the threads fetcher.join() writer.join() except KeyboardInterrupt: print("Finishing threads due to keyboard interruption.") fetcher.join() writer.join() except Exception as e: print("Error encountered: ", e) finally: if conn: conn.close() def get_rt_db_path(db_path: str, db_extension: str = '.RT'): db_dir, db_file = os.path.split(db_path) db_name, _ = os.path.splitext(db_file) if db_name.endswith('_Raw'): db_name = db_name[:-4] rt_db_name = db_name + '_' + db_extension return os.path.join(db_dir, rt_db_name) def main(): db_path = input("Enter absolute path of raw.db file: ") try: maxsize=1000 db_writer = Queue(maxsize=maxsize) tables = Tables(max_buffer_length=maxsize) rt_db_path = get_rt_db_path(db_path) db_writer_thread = Thread(target=tables.process_db_writes, args=(db_writer, rt_db_path)) # Start the db_writer_thread db_writer_thread.start() fetch_raw_data(db_path, rt_db_path, maxsize, db_writer) db_writer_thread.join() except KeyboardInterrupt: db_writer_thread.join() except Exception: db_writer_thread.join()
核心问题原因
列表引用复用导致数据覆盖
在insert_batch_data中,将batch放入队列后立即调用batch.clear()。由于列表是可变对象,队列中存储的是列表引用而非副本,后续的清空操作会直接修改队列内的列表内容,导致消费线程读到空数据。不合理的断言阻断剩余数据写入
process_db_writes和add_to_writer_queue中的断言要求所有数据长度必须等于max_buffer_length,但最后一批剩余数据长度必然小于该值,触发断言错误后程序中断,剩余数据无法写入。线程参数传递格式错误
fetch_raw_data中启动线程时,args=(buffer)未加逗号,会被解析为单个元素的序列而非元组,导致线程函数接收参数异常。缺少终止信号导致线程阻塞
populate_queue完成数据生成后未向队列发送None终止信号,insert_batch_data中的循环会一直阻塞等待新数据,无法正常退出。
修复步骤
1. 传递列表副本避免引用复用
修改insert_batch_data,放入队列时创建列表副本,防止后续操作修改队列内的数据:
def insert_batch_data(self, buffer: Queue, table_name: str, writer_queue: Queue): query = self.get_query(table_name) batch = [] while True: if buffer.empty(): time.sleep(0.1) # 缩短等待时间提升响应性 continue item = buffer.get() if item is None: print("Reached Sentinal Value... Exiting thread...") break batch.append(item) if len(batch) == self.max_buffer_length: # 传递列表副本而非引用 self.add_to_writer_queue(query, table_name, batch.copy(), writer_queue) print(f"Number of items added to writer_queue: {len(batch)}") batch.clear() if batch: self.add_to_writer_queue(query, table_name, batch.copy(), writer_queue)
2. 移除不合理断言,兼容剩余数据
修改process_db_writes和add_to_writer_queue,仅对满批次数据做校验:
def process_db_writes(self, writer_queue: Queue, db_path: str): try: conn = sqlite3.connect(db_path) cursor = conn.cursor() while True: while writer_queue.empty(): time.sleep(0.1) query, table_name, data = writer_queue.get() self.insert_records(cursor, conn, query, data, table_name) writer_queue.task_done() # 标记任务完成,避免队列阻塞 except Exception as e: print(f"Error encountered in process_db_writes: {str(e)}") finally: cursor.close() conn.close() def add_to_writer_queue(self, query: str, table_name: str, batch: List, writer_queue: Queue): while writer_queue.full(): time.sleep(0.1) # 仅对满批次数据做断言 if len(batch) == self.max_buffer_length: assert len(batch) == self.max_buffer_length, f"Expected {self.max_buffer_length} items, received: {len(batch)}\n" writer_queue.put((query, table_name, batch))
3. 修复线程参数传递错误
调整fetch_raw_data中的线程启动代码,确保参数为元组格式:
def fetch_raw_data(db_path: str, rt_db_path: str, db_writer: Queue, max_buffer_length: int, ): try: conn = sqlite3.connect(db_path) buffer = Queue(maxsize=max_buffer_length) raw = Raw(max_buffer_length=max_buffer_length) # 添加逗号确保参数是元组 fetcher = Thread(target=raw.populate_queue, args=(buffer,)) writer = Thread(target=raw.insert_data, args=(buffer, rt_db_path, db_writer)) fetcher.start() writer.start() fetcher.join() buffer.put(None) # 发送终止信号 writer.join() except KeyboardInterrupt: print("Finishing threads due to keyboard interruption.") buffer.put(None) fetcher.join() writer.join() except Exception as e: print("Error encountered: ", e) finally: if conn: conn.close()
4. 添加终止信号确保线程正常退出
修改populate_queue,完成数据生成后发送终止信号:
def populate_queue(buffer: Queue): for i in range(1_000_000): while buffer.full(): time.sleep(0.1) buffer.put(RawData(time=i)) buffer.put(None) # 发送终止信号
5. 修复main函数参数顺序错误
调整main中fetch_raw_data的参数顺序,匹配函数定义:
def main(): db_path = input("Enter absolute path of raw.db file: ") try: maxsize=1000 db_writer = Queue(maxsize=maxsize) tables = Tables(max_buffer_length=maxsize) rt_db_path = get_rt_db_path(db_path) db_writer_thread = Thread(target=tables.process_db_writes, args=(db_writer, rt_db_path)) db_writer_thread.start() # 调整参数顺序 fetch_raw_data(db_path, rt_db_path, db_writer, maxsize) db_writer.join() db_writer_thread.join() except KeyboardInterrupt: db_writer_thread.join() except Exception as e: print(f"Main error: {e}") db_writer_thread.join()
额外优化建议
- 缩短等待时间:将
time.sleep(5)、time.sleep(2)改为time.sleep(0.1),提升线程响应速度。 - 完善任务管理:使用
Queue.task_done()和Queue.join()确保所有数据处理完成后再退出程序。 - 增强日志:在关键步骤添加详细日志,便于定位问题。
内容的提问来源于stack exchange,提问作者mraabhijit

