文章
duckdb 数据治理新工具
今天无意了解到duckdb,可以方便读取本地csv、pandas DataFrame、json等数据然后直接使用sql进行处理,体验了一番,确实不错
安装
pip install duckdbdemo01
import duckdb
print("读取单个指定csv文件")
# 从文件读取
res = duckdb.read_csv("input/data.csv")
print(res)
print("=" * 50)
print("从文件夹读取多个csv文件")
# 从文件夹读取
res = duckdb.read_csv("input/*.csv")
print(res)
print("=" * 50)
# 指定内部格式 dtype 可以指定列类型
res = duckdb.read_csv("input/data.csv", header=False, sep=",", dtype=["varchar", "int", "int"])
print(res)
print("=" * 50)
print("在sql中读取文件")
res = duckdb.sql("SELECT * FROM 'input/data.csv'")
print(res)
print("=" * 50)
print("在sql内部调用read_csv")
res = duckdb.sql("SELECT * FROM read_csv('input/data.csv')")
print(res)demo02
import duckdb
res = duckdb.read_json("input/data.json")
print(res)
res = duckdb.read_json("input/*.json")
print(res)
res = duckdb.sql("SELECT * FROM 'input/data.json'")
print(res)
res = duckdb.sql("SELECT * FROM read_json_auto('input/*.json')")
print(res)demo03
import duckdb
import pandas as pd
test_df = pd.DataFrame.from_dict({"i": [1, 2, 3, 4], "j": ["one", "two", "three", "four"]})
print(test_df)
res = duckdb.sql("SELECT * FROM test_df")
print(res)
print(res.fetchall())demo04
import duckdb
import pandas as pd
my_dictionary = {}
my_dictionary["test_df"] = pd.DataFrame.from_dict({"i": [1, 2, 3, 4], "j": ["one", "two", "three", "four"]})
duckdb.register("test_df_view", my_dictionary["test_df"])
print(duckdb.sql("SELECT * FROM test_df_view").fetchall())demo05
import duckdb
from duckdb.sqltypes import VARCHAR
from faker import Faker
def generate_random_name():
fake = Faker()
return fake.name()
duckdb.create_function("random_name", generate_random_name, [], VARCHAR)
res = duckdb.sql("SELECT random_name()").fetchall()
print(res)
数据治理应用:
生成脏数据:
import csv
import random
from datetime import datetime, timedelta
# 固定随机种子,保证每次生成的数据一致(方便对照)
random.seed(42)
# 定义数据维度
sources = ['web_crawl', 'pdf_parse', 'internal_doc', 'api_import']
# 用于生成正常文本的模板
templates = [
"大模型微调语料:{keyword} 在 {date} 发生了重要变化,影响范围包括 {scope}。",
"合同条款摘要:甲方 {party_a} 与乙方 {party_b} 就 {project} 达成合作意向。",
"医疗诊断记录:患者主诉 {symptom},初步诊断为 {diagnosis},建议 {suggestion}。",
"技术文档描述:系统 {system} 在版本 {version} 中修复了 {bug} 漏洞,提升了性能。",
"财报数据摘要:{company} 本季度营收 {revenue} 亿元,同比增长 {growth}%。"
]
keywords = ['AI', '数据安全', '云计算', '边缘计算']
scopes = ['华东区', '华南区', '全国']
party_a = ['华为', '腾讯', '阿里', '字节跳动']
party_b = ['中国移动', '中国电信', '中国联通']
projects = ['智慧城市', '边缘计算节点', '大模型训练']
symptoms = ['头痛', '发热', '关节痛', '无症状']
diagnosis = ['普通感冒', '流感', '高血压', '健康']
suggestions = ['多休息', 'CT检查', '血常规', '随访观察']
systems = ['Kubernetes', 'MySQL', 'Redis', 'Nginx']
versions = ['1.0.0', '2.1.3', '3.0.5']
bugs = ['内存泄漏', '死锁', '安全注入', '缓存穿透']
companies = ['比亚迪', '宁德时代', '理想汽车', '小米']
revenues = [100.5, 250.8, 320.1, 88.6]
growths = [15.6, 22.3, 8.9, 35.1]
# 预置一些完全相同的重复文本(用于去重测试)
duplicate_texts = [
"这是完全重复的文本A,用于测试去重逻辑。",
"这是完全重复的文本B,用于测试去重逻辑。"
]
# 打开 CSV 文件写入
with open('input/sample_data.csv', 'w', newline='', encoding='utf-8-sig') as f:
writer = csv.writer(f)
# 写入表头
writer.writerow(['id', 'source', 'raw_text', 'quality_score', 'ingest_date'])
for i in range(1, 101):
source = random.choice(sources)
rand_val = random.random()
# ---------- 构造脏数据 ----------
if rand_val < 0.10:
# 10%:完全解析失败 (模拟 PDF 损坏)
raw_text = None
quality_score = None
elif rand_val < 0.20:
# 10%:空字符串或全空格 (模拟爬虫抓取空页面)
raw_text = "" if random.random() > 0.5 else " \t "
quality_score = 0.0 # 给个低分
elif rand_val < 0.30:
# 10%:故意重复的文本 (模拟重复抓取)
raw_text = random.choice(duplicate_texts)
quality_score = round(random.uniform(0.6, 0.95), 2)
else:
# 70%:正常文本(但偶尔带换行符和多余空格)
template = random.choice(templates)
raw_text = template.format(
keyword=random.choice(keywords),
date=(datetime.now() - timedelta(days=random.randint(1, 365))).strftime("%Y-%m-%d"),
scope=random.choice(scopes),
party_a=random.choice(party_a),
party_b=random.choice(party_b),
project=random.choice(projects),
symptom=random.choice(symptoms),
diagnosis=random.choice(diagnosis),
suggestion=random.choice(suggestions),
system=random.choice(systems),
version=random.choice(versions),
bug=random.choice(bugs),
company=random.choice(companies),
revenue=random.choice(revenues),
growth=random.choice(growths)
)
# 20% 的概率插入换行符(模拟真实文档的段落)
if random.random() > 0.8:
raw_text = raw_text.replace("。", "。\n")
# 20% 的概率加多余空格
if random.random() > 0.8:
raw_text = raw_text.replace(" ", " ")
quality_score = round(random.uniform(0.5, 0.99), 2)
# 如果文本为空或 None,质量分强制为空
if raw_text in [None, "", " \t "]:
quality_score = None
# 随机日期(最近30天)
ingest_date = (datetime.now() - timedelta(days=random.randint(0, 29))).strftime('%Y-%m-%d')
# 写入一行
writer.writerow([i, source, raw_text, quality_score, ingest_date])
print("✅ 数据生成成功!文件路径:./sample_data.csv")
print("📊 共生成 100 条记录,包含脏数据、重复值、缺失值。")治理脚本:
import duckdb
# 1. 直接查询 CSV 文件(零内存加载,流式读取)
print("===== 原始数据预览(前5行) =====")
preview = duckdb.sql(
"""
SELECT *
FROM 'sample_data.csv' LIMIT 5
"""
).df()
print(preview)
# 2. 执行数据质量检查(看看有多少脏数据)
print("\n===== 数据质量报告 =====")
quality_report = duckdb.sql(r"""
SELECT
COUNT(*) AS total_rows,
SUM(CASE WHEN raw_text IS NULL THEN 1 ELSE 0 END) AS null_count,
SUM(CASE WHEN REGEXP_REPLACE(raw_text, '\s', '', 'g') = '' THEN 1 ELSE 0 END) AS empty_count,
SUM(CASE WHEN quality_score IS NULL THEN 1 ELSE 0 END) AS null_score_count,
COUNT(DISTINCT raw_text) AS unique_texts
FROM 'sample_data.csv'
""").df()
print(quality_report)
# 3. 清洗 找出重复文本
print("\n===== 重复文本统计(出现次数 > 1) =====")
duplicates = duckdb.sql(
"""
SELECT raw_text,
count(1) as cnt,
STRING_AGG(CAST(id AS VARCHAR), ', ') AS ids
FROM 'sample_data.csv'
GROUP BY raw_text
HAVING count (1) > 1
"""
).df()
print(duplicates)
res = duckdb.sql(
r"""
WITH cleaned AS (SELECT id,
source,
-- 1. 清洗文本:不再用 TRIM,直接用正则把 \t 多余空格全部处理掉
REGEXP_REPLACE(
COALESCE(raw_text, ''),
'\s+',
' ',
'g'
) AS clean_text, quality_score,
ingest_date,
-- 2. 计算有效长度:去除所有空白符后还剩几个可见字符
LENGTH(
REGEXP_REPLACE(
COALESCE(raw_text, ''),
'\s',
'',
'g'
)
) AS char_len
FROM 'sample_data.csv'
WHERE raw_text IS NOT NULL
-- 3. 关键修改:判断“去除所有空白后还有没有可见字符”
AND LENGTH (REGEXP_REPLACE(raw_text
, '\s'
, ''
, 'g'))
> 0
)
,
-- 后续 deduped 和 ranked 逻辑完全不用动
deduped AS (
SELECT
*, ROW_NUMBER() OVER (PARTITION BY clean_text ORDER BY quality_score DESC NULLS LAST) AS rn
FROM cleaned
WHERE char_len >= 10
)
, ranked AS (
SELECT
*, ROW_NUMBER() OVER (PARTITION BY source ORDER BY quality_score DESC NULLS LAST) AS source_rank
FROM deduped
WHERE rn = 1
)
SELECT source,
clean_text,
quality_score,
char_len,
ingest_date
FROM ranked
WHERE source_rank <= 3
ORDER BY source, quality_score DESC;
"""
)
print(res)附录: