如何使用DuckDB批量更新Parquet文件字段值及替换字符
用DuckDB修改Parquet文件内容
一、将FirstName字段的"John"替换为"Alex"
DuckDB支持递归读取嵌套目录的Parquet文件,结合数据修改和导出操作可实现需求,步骤如下:
读取所有Parquet文件
使用read_parquet开启recursive=true参数,一次性加载/mainfolder下所有层级的Parquet文件,同时保留原文件路径用于后续还原目录结构:CREATE TABLE temp_data AS SELECT *, filename AS source_file FROM read_parquet('/mainfolder/**/*.parquet', recursive=true, filename=true);修改目标字段
用CASE语句精确替换FirstName字段的"John"为"Alex"(若需替换所有包含"John"的内容,改用replace(FirstName, 'John', 'Alex')):CREATE TABLE modified_data AS SELECT CASE WHEN FirstName = 'John' THEN 'Alex' ELSE FirstName END AS FirstName, * EXCLUDE FirstName, source_file FROM temp_data;导出修改后的数据
- 无需保留原目录结构时,直接导出到新目录:
COPY modified_data EXCLUDE source_file TO '/new_mainfolder' (FORMAT PARQUET); - 需要保留原目录结构时,可结合Python脚本拆分路径并写入对应文件夹:
import duckdb import os conn = duckdb.connect() conn.execute(""" CREATE TABLE modified_data AS SELECT CASE WHEN FirstName = 'John' THEN 'Alex' ELSE FirstName END AS FirstName, * EXCLUDE FirstName, filename AS source_file FROM read_parquet('/mainfolder/**/*.parquet', recursive=true, filename=true); """) # 获取所有唯一文件路径 paths = conn.execute("SELECT DISTINCT source_file FROM modified_data").fetchall() for path in paths: source_path = path[0] target_path = source_path.replace('/mainfolder/', '/new_mainfolder/') os.makedirs(os.path.dirname(target_path), exist_ok=True) conn.execute(f""" COPY (SELECT * EXCLUDE source_file FROM modified_data WHERE source_file = '{source_path}') TO '{target_path}' (FORMAT PARQUET); """)
- 无需保留原目录结构时,直接导出到新目录:
注意:DuckDB无法直接覆盖原Parquet文件,建议先导出到新目录,确认无误后再替换原文件。
二、批量替换所有字段中的"e"为"f"
DuckDB没有原生的批量遍历所有字段替换功能,因为字段可能包含数值、日期等非字符串类型,这类字段无法执行字符串替换。需先筛选出字符串类型字段,再逐个处理:
方法1:纯SQL手动列举字段
字段数量不多时,直接写出所有字符串字段的替换逻辑:
CREATE TABLE modified_all AS SELECT replace(FirstName, 'e', 'f') AS FirstName, replace(LastName, 'e', 'f') AS LastName, -- 非字符串字段直接保留 Age, CreateTime FROM read_parquet('/mainfolder/**/*.parquet', recursive=true);
方法2:结合Python动态生成SQL
字段较多时,用Python自动获取字符串类型字段并生成替换语句:
import duckdb conn = duckdb.connect() # 读取数据到临时表 conn.execute("CREATE TABLE temp_data AS SELECT * FROM read_parquet('/mainfolder/**/*.parquet', recursive=true);") # 查询所有字符串类型字段 string_columns = conn.execute(""" SELECT column_name FROM information_schema.columns WHERE table_name = 'temp_data' AND data_type IN ('VARCHAR', 'STRING'); """).fetchall() # 生成替换字段的SQL片段 replace_clauses = [f"replace({col[0]}, 'e', 'f') AS {col[0]}" for col in string_columns] # 生成保留非字符串字段的SQL片段 other_columns = conn.execute(""" SELECT column_name FROM information_schema.columns WHERE table_name = 'temp_data' AND data_type NOT IN ('VARCHAR', 'STRING'); """).fetchall() other_clauses = [col[0] for col in other_columns] # 拼接并执行完整SQL full_sql = f"CREATE TABLE modified_all AS SELECT {', '.join(replace_clauses + other_clauses)} FROM temp_data;" conn.execute(full_sql) # 导出数据 conn.execute("COPY modified_all TO '/new_mainfolder_all' (FORMAT PARQUET);")
内容的提问来源于stack exchange,提问作者sojim2
相关产品推荐
相关产品推荐

