Многопоточные операции в Python
Эта страница демонстрирует, как одновременно вставлять данные и читать их из базы данных DuckDB в нескольких потоках Python. Это может быть полезно в сценариях, где поступают новые данные, и анализ должен периодически пересчитываться. Обратите внимание, что всё это происходит в рамках одного процесса Python (подробности о параллельности DuckDB см. в FAQ). Вы можете следовать за примером в этом блокноте Google Colab.
Настройка
Сначала импортируйте DuckDB и несколько модулей из стандартной библиотеки Python. Примечание: если вы используете Pandas, добавьте import pandas в начало скрипта (так как он должен быть импортирован до многопоточных операций). Затем подключитесь к файловой базе данных DuckDB и создайте пример таблицы для хранения вставленных данных. Эта таблица будет отслеживать имя потока, который выполнил вставку, и автоматически вставлять метку времени, когда эта вставка произошла, используя DEFAULT выражение.
import duckdb
from threading import Thread, current_thread
import random
duckdb_con = duckdb.connect('my_peristent_db.duckdb')
# Use connect without parameters for an in-memory database
# duckdb_con = duckdb.connect()
duckdb_con.execute("""
CREATE OR REPLACE TABLE my_inserts (
thread_name VARCHAR,
insert_time TIMESTAMP DEFAULT current_timestamp
)
""") Функции чтения и записи
Далее, определите функции, которые будут выполняться потоками записи и чтения. Каждый поток должен использовать метод .cursor() для создания локального подключения к тому же файлу DuckDB на основе исходного подключения. Этот подход также работает с базами данных DuckDB в памяти.
def write_from_thread(duckdb_con):
# Create a DuckDB connection specifically for this thread
local_con = duckdb_con.cursor()
# Insert a row with the name of the thread. insert_time is auto-generated.
thread_name = str(current_thread().name)
result = local_con.execute("""
INSERT INTO my_inserts (thread_name)
VALUES (?)
""", (thread_name,)).fetchall()
def read_from_thread(duckdb_con):
# Create a DuckDB connection specifically for this thread
local_con = duckdb_con.cursor()
# Query the current row count
thread_name = str(current_thread().name)
results = local_con.execute("""
SELECT
? AS thread_name,
count(*) AS row_counter,
current_timestamp
FROM my_inserts
""", (thread_name,)).fetchall()
print(results) Создание потоков
Определите количество потоков записи и чтения и список, чтобы отслеживать все созданные потоки. Затем создайте сначала потоки записи, а затем чтения. После этого перемешайте их, чтобы они запускались в случайном порядке, что симулирует одновременные операции записи и чтения. Обратите внимание, что потоки пока не были выполнены, только определены.
write_thread_count = 50
read_thread_count = 5
threads = []
# Create multiple writer and reader threads (in the same process)
# Pass in the same connection as an argument
for i in range(write_thread_count):
threads.append(Thread(target = write_from_thread,
args = (duckdb_con,),
name = 'write_thread_' + str(i)))
for j in range(read_thread_count):
threads.append(Thread(target = read_from_thread,
args = (duckdb_con,),
name = 'read_thread_' + str(j)))
# Shuffle the threads to simulate a mix of readers and writers
random.seed(6) # Set the seed to ensure consistent results when testing
random.shuffle(threads) Запуск потоков и отображение результатов
Теперь запустите все потоки параллельно, затем подождите, пока все они завершат работу, прежде чем выводить результаты. Обратите внимание, что метки времени чтений и записей чередуются, как ожидалось, из-за рандомизации.
# Kick off all threads in parallel
for thread in threads:
thread.start()
# Ensure all threads complete before printing final results
for thread in threads:
thread.join()
print(duckdb_con.execute("""
SELECT *
FROM my_inserts
ORDER BY
insert_time
""").df())
© Copyright 2018–2024 Stichting DuckDB Foundation
Licensed under the MIT License.
https://duckdb.org/docs/guides/python/multiple_threads.html