You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在DuckDB中实现无重复增量写入及动态添加列?

解决方案

针对你的两个需求,我们可以通过DuckDB的元数据查询、动态SQL以及数据比对来实现,以下是修改后的完整代码和关键逻辑说明:

核心逻辑梳理

  1. 表存在性检查:通过查询information_schema.tables判断目标表是否已存在
  2. 避免重复行:利用EXCEPT子句过滤已存在的行,或通过唯一键使用ON CONFLICT DO NOTHING(推荐后者,效率更高)
  3. 新增列兼容:读取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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 07:20:27