여러 Python 스레드

여러 Python 스레드 (Multiple Python Threads)

이 페이지는 여러 Python 스레드에서 동시에 DuckDB 데이터베이스에 삽입하면서 읽는 방법을 보여줘요. 새 데이터가 계속 유입되고 분석을 주기적으로 다시 실행해야 하는 시나리오에서 유용할 수 있어요. 참고로 이 모든 작업은 단일 Python 프로세스 안에서 일어납니다(DuckDB 동시성에 대한 자세한 내용은 FAQ 참고). 이 Google Colab 노트북을 따라 해 볼 수도 있어요.

설정 (Setup)

먼저 DuckDB와 Python 표준 라이브러리의 여러 모듈을 임포트합니다. 참고: Pandas를 쓴다면 스크립트 맨 위에 import pandas를 추가하세요. 멀티스레딩 이전에 먼저 임포트돼야 하기 때문이에요. 그다음 파일 백업(file-backed) DuckDB 데이터베이스에 연결하고, 삽입된 데이터를 저장할 예시 테이블을 만듭니다.

이 테이블은 삽입을 완료한 스레드의 이름을 기록하고, DEFAULT 표현식으로 삽입 시점의 타임스탬프를 자동으로 넣어요.

import duckdb
from threading import Thread, current_thread
import random

duckdb_con = duckdb.connect('my_persistent_db.duckdb')
# 인메모리 데이터베이스는 파라미터 없이 connect()를 씁니다
# duckdb_con = duckdb.connect()
duckdb_con.execute("""
    CREATE OR REPLACE TABLE my_inserts (
        thread_name VARCHAR,
        insert_time TIMESTAMP DEFAULT current_timestamp
    )
""")

읽기·쓰기 함수 (Reader and Writer Functions)

다음으로 쓰기 스레드와 읽기 스레드가 실행할 함수를 정의해요. 각 스레드는 원본 연결을 바탕으로 같은 DuckDB 파일에 대한 스레드 로컬 연결을 만들기 위해 반드시 .cursor() 메서드를 사용해야 합니다. 이 방식은 인메모리 DuckDB 데이터베이스에서도 동일하게 동작해요.

def write_from_thread(duckdb_con):
    # 이 스레드만을 위한 DuckDB 연결 생성
    local_con = duckdb_con.cursor()
    # 스레드 이름을 넣은 행 삽입. insert_time은 자동 생성됨.
    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):
    # 이 스레드만을 위한 DuckDB 연결 생성
    local_con = duckdb_con.cursor()
    # 현재 행 개수 조회
    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)

스레드 생성 (Create Threads)

쓸 쓰기·읽기 스레드 개수를 정하고, 만들 스레드를 모두 추적할 리스트를 정의합니다. 그다음 먼저 쓰기 스레드를, 다음으로 읽기 스레드를 만듭니다. 이어서 셔플해서 무작위 순서로 시작되게 하여 읽기와 쓰기가 동시에 일어나는 상황을 흉내 냅니다. 이 단계에서는 스레드를 정의만 했지 아직 실행하지는 않았어요.

write_thread_count = 50
read_thread_count = 5
threads = []

# (같은 프로세스 안에서) 여러 쓰기·읽기 스레드 생성
# 같은 연결을 인자로 전달
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)))

# 읽기와 쓰기가 섞이도록 스레드 셔플
random.seed(6) # 테스트 결과를 일관되게 만들기 위해 시드 설정
random.shuffle(threads)

스레드 실행과 결과 확인 (Run Threads and Show Results)

이제 모든 스레드를 병렬로 시작하고, 결과를 출력하기 전에 전부 끝날 때까지 기다립니다. 무작위화 덕분에 읽기·쓰기 스레드의 타임스탬프가 예상대로 섞여 나타나는 걸 볼 수 있어요.

# 모든 스레드를 병렬로 시작
for thread in threads:
    thread.start()

# 최종 결과 출력 전에 모든 스레드 완료 보장
for thread in threads:
    thread.join()

print(duckdb_con.execute("""
    SELECT *
    FROM my_inserts
    ORDER BY
        insert_time
""").df())

더 알아보기 (Learn more)