prevent race condiction on drop_and_create_table_if_schema_mismatch
This commit is contained in:
parent
0253b15a2c
commit
b1e7a8b288
1 changed files with 27 additions and 22 deletions
|
|
@ -1,6 +1,7 @@
|
||||||
from typing import TYPE_CHECKING, Any, Dict, Optional, Type, Union
|
from typing import TYPE_CHECKING, Any, Dict, Optional, Type, Union
|
||||||
|
|
||||||
import duckdb
|
import duckdb
|
||||||
|
import threading
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
|
@ -13,6 +14,9 @@ if TYPE_CHECKING:
|
||||||
|
|
||||||
INDEX_KEY = "index"
|
INDEX_KEY = "index"
|
||||||
|
|
||||||
|
# Lock to prevent multiple threads from creating the same table at the same time
|
||||||
|
drop_create_table_lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
def get_table_schema_as_dict(conn: duckdb.DuckDBPyConnection, table_name: str) -> dict:
|
def get_table_schema_as_dict(conn: duckdb.DuckDBPyConnection, table_name: str) -> dict:
|
||||||
result = conn.execute(f"PRAGMA table_info('{table_name}')").fetchall()
|
result = conn.execute(f"PRAGMA table_info('{table_name}')").fetchall()
|
||||||
|
|
@ -48,6 +52,7 @@ def model_to_sql_column_definitions(model: Type[BaseModel]) -> dict:
|
||||||
|
|
||||||
|
|
||||||
def drop_and_create_table_if_schema_mismatch(db_path: str, table_name: str, model: Type[BaseModel]):
|
def drop_and_create_table_if_schema_mismatch(db_path: str, table_name: str, model: Type[BaseModel]):
|
||||||
|
with drop_create_table_lock:
|
||||||
with duckdb.connect(db_path) as conn:
|
with duckdb.connect(db_path) as conn:
|
||||||
# Get the current schema from the database
|
# Get the current schema from the database
|
||||||
try:
|
try:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue