如何在DuckDB中实现无重复增量写入及动态添加列?
解决方案
针对你的两个需求,我们可以通过DuckDB的元数据查询、动态SQL以及数据比对来实现,以下是修改后的完整代码和关键逻辑说明:
核心逻辑梳理
- 表存在性检查:通过查询
information_schema.tables判断目标表是否已存在 - 避免重复行:利用
EXCEPT子句过滤已存在的行,或通过唯一键使用ON CONFLICT DO NOTHING(推荐后者,效率更高) - 新增列兼容:读取CSV的列结构,对比现有表列,自动添加缺失的新列
修改后的完整代码
import time from watchdog.observers import Observer from watchdog.events import PatternMatchingEventHandler import duckdb path = "landing/persistent/" con = duckdb.connect(database='formatted/my-db.duckdb', read_only=False) class EventHandler(PatternMatchingEventHandler): def __init__(self): super().__init__( patterns=["*.csv"], ignore_directories=True, case_sensitive=False, ) def _table_exists(self, cursor, table_name): # 检查目标表是否存在 cursor.execute(""" SELECT EXISTS( SELECT 1 FROM information_schema.tables WHERE table_name = ? AND table_schema = 'main' ) """, [table_name]) return cursor.fetchone()[0] def _get_table_columns(self, cursor, table_name): # 获取现有表的列名集合 cursor.execute(""" SELECT column_name FROM information_schema.columns WHERE table_name = ? AND table_schema = 'main' ORDER BY ordinal_position """, [table_name]) return {row[0] for row in cursor.fetchall()} def _get_csv_columns(self, cursor, csv_path): # 获取CSV文件的列名集合和列类型(仅读取结构,不加载数据) cursor.execute("SELECT * FROM read_csv_auto(?, header=True, n_rows=0)", [csv_path]) result = cursor.description return {col[0] for col in result}, {col[0]: col[1] for col in result} def on_created(self, event): filename = event.src_path.split("/")[-1] table_name = filename.split(".")[0] csv_path = event.src_path try: cursor = con.cursor() if not self._table_exists(cursor, table_name): # 表不存在,直接创建新表 cursor.execute(f"CREATE TABLE {table_name} AS SELECT * FROM read_csv_auto(?)", [csv_path]) print(f"新表 {table_name} 创建并导入完成") else: # 表已存在,先处理新增列 existing_cols = self._get_table_columns(cursor, table_name) csv_cols, csv_col_types = self._get_csv_columns(cursor, csv_path) # 添加缺失的新列 new_cols = csv_cols - existing_cols for col in new_cols: col_type = csv_col_types[col] cursor.execute(f"ALTER TABLE {table_name} ADD COLUMN {col} {col_type}") print(f"为表 {table_name} 添加新列: {col} ({col_type})") # 导入数据并避免重复行(二选一) # 方式1:有唯一键时用(性能更高) # cursor.execute(f""" # INSERT INTO {table_name} # SELECT * FROM read_csv_auto(?) # ON CONFLICT (id) DO NOTHING # """, [csv_path]) # 方式2:无唯一键时用(全量比对) cursor.execute(f""" INSERT INTO {table_name} SELECT * FROM read_csv_auto(?) EXCEPT SELECT * FROM {table_name} """, [csv_path]) print(f"表 {table_name} 已更新,导入了非重复数据") # 打印当前所有表 cursor.execute("show tables") print("当前数据库表列表:", cursor.fetchall()) except Exception as e: print(f"处理文件 {filename} 出错: {str(e)}") finally: cursor.close() event_handler = EventHandler() observer = Observer() observer.schedule(event_handler, path, recursive=True) observer.start() try: while True: time.sleep(1) except KeyboardInterrupt: observer.stop() observer.join()
关键细节说明
1. 避免重复行的两种方式
- 方式1(推荐):唯一键冲突处理
如果你的CSV数据有唯一标识列(比如id),使用ON CONFLICT (唯一键列) DO NOTHING可以高效跳过已存在的行,性能远优于全量比对。 - 方式2:无唯一键时的全量比对
用EXCEPT子句对比CSV数据和现有表数据,只插入不存在的行。注意这种方式会对全表数据进行比对,数据量大时性能会下降,建议尽量定义唯一键。
2. 新增列自动处理
- 通过
read_csv_auto(..., n_rows=0)读取CSV的列结构(不加载实际数据),获取列名和类型 - 对比现有表的列,自动执行
ALTER TABLE ADD COLUMN添加缺失的列 - DuckDB会自动处理新列的NULL值(导入时CSV中没有的旧数据行,新列值为NULL)
3. 其他优化
- 移除了不必要的全局变量声明,代码结构更简洁
- 增加了详细的打印日志,方便排查问题
- 使用
super()简化了父类初始化逻辑
内容的提问来源于stack exchange,提问作者Norhther
相关产品推荐
相关产品推荐

