Zum Inhalt springen

Mehrere Python-Threads

Diese Seite zeigt, wie Sie gleichzeitig aus mehreren Python-Threads in eine DuckDB-Datenbank einfügen und daraus lesen. Das kann nützlich sein, wenn neue Daten eintreffen und eine Analyse regelmäßig neu ausgeführt werden soll. Beachten Sie, dass alles innerhalb eines einzelnen Python-Prozesses stattfindet (siehe die FAQ für Details zur DuckDB-Nebenläufigkeit). Sie können das Beispiel in diesem Google-Colab-Notebook nachvollziehen.

Einrichtung

Importieren Sie zuerst DuckDB und mehrere Module aus der Python-Standardbibliothek. Hinweis: Wenn Sie Pandas verwenden, fügen Sie import pandas ebenfalls am Anfang des Skripts hinzu (es muss vor dem Multi-Threading importiert werden). Stellen Sie dann eine Verbindung zu einer dateibasierten DuckDB-Datenbank her und legen Sie eine Beispieltabelle an, in der die eingefügten Daten gespeichert werden. Diese Tabelle merkt sich den Namen des Threads, der den Insert ausgeführt hat, und setzt automatisch den Zeitstempel des Inserts über den DEFAULT-Ausdruck.

import duckdb
from threading import Thread, current_thread
import random
duckdb_con = duckdb.connect('my_persistent_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
)
""")

Reader- und Writer-Funktionen

Als Nächstes definieren Sie Funktionen, die von den Writer- und Reader-Threads ausgeführt werden. Jeder Thread muss die Methode .cursor() verwenden, um eine thread-lokale Verbindung zur selben DuckDB-Datei auf Basis der ursprünglichen Verbindung zu erzeugen. Dieser Ansatz funktioniert auch mit In-Memory-DuckDB-Datenbanken.

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)

Threads erzeugen

Wir legen fest, wie viele Writer und Reader verwendet werden, und definieren eine Liste, in der alle erzeugten Threads festgehalten werden. Anschließend erzeugen wir zuerst Writer- und dann Reader-Threads. Danach mischen wir sie, sodass sie in zufälliger Reihenfolge gestartet werden und gleichzeitige Writer und Reader simuliert werden. Beachten Sie, dass die Threads noch nicht ausgeführt, sondern nur definiert wurden.

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)

Threads ausführen und Ergebnisse anzeigen

Starten Sie nun alle Threads parallel und warten Sie, bis sie fertig sind, bevor Sie die Ergebnisse ausgeben. Beachten Sie, dass die Zeitstempel von Readern und Writern wie erwartet durcheinanderliegen, weil die Reihenfolge zufällig ist.

# 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())