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

向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()
问题分析与修复方案

核心问题原因

  1. 列表引用复用导致数据覆盖
    在insert_batch_data中,将batch放入队列后立即调用batch.clear()。由于列表是可变对象,队列中存储的是列表引用而非副本,后续的清空操作会直接修改队列内的列表内容,导致消费线程读到空数据。

  2. 不合理的断言阻断剩余数据写入
    process_db_writes和add_to_writer_queue中的断言要求所有数据长度必须等于max_buffer_length,但最后一批剩余数据长度必然小于该值,触发断言错误后程序中断,剩余数据无法写入。

  3. 线程参数传递格式错误
    fetch_raw_data中启动线程时,args=(buffer)未加逗号,会被解析为单个元素的序列而非元组,导致线程函数接收参数异常。

  4. 缺少终止信号导致线程阻塞
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 11:44:51