Article
第1章:DuckDB 入门概览
初识 DuckDB,了解其定位、特点与适用场景。
1.1 什么是 DuckDB
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| DuckDB | 嵌入式分析型数据库管理系统(Analytical DBMS),专为 OLAP 工作负载设计,支持 SQL,以列式存储和向量化执行引擎为核心。 | 不适用于高并发 OLTP 场景(如 Web 应用后端写密集操作)。 |
| 嵌入式数据库 | 直接嵌入到应用程序中运行,无需独立服务器进程,类似 SQLite。 | 数据库文件本地化,适合单机或边缘计算场景。 |
| 分析型数据库(OLAP) | 优化用于复杂查询、聚合分析、大规模数据扫描,而非频繁的增删改操作。 | 查询性能远高于传统行式数据库在分析任务中的表现。 |
| 开源协议 | 使用 MIT 许可证,允许商业使用、修改和分发。 | 可自由集成至闭源项目中。 |
1.2 DuckDB 与 SQLite、PostgreSQL、Pandas 的对比
| 对比维度 | DuckDB | SQLite | PostgreSQL | Pandas |
|---|---|---|---|---|
| 类型 | 嵌入式分析型数据库(OLAP) | 嵌入式事务型数据库(OLTP) | 客户端-服务器关系型数据库(通用) | 内存数据处理库(Python) |
| 存储模式 | 列式存储为主,优化聚合查询 | 行式存储 | 行式存储(支持列存扩展) | 内存中按列组织(DataFrame) |
| 执行引擎 | 向量化执行,CPU 缓存友好 | 解释执行 | 解释/ JIT 编译 | 向量化函数(NumPy 底层) |
| 并发支持 | 单写多读(MVCC),轻量级并发 | 单连接写,其余只读 | 高并发多用户支持 | GIL 限制,非线程安全 |
| 使用场景 | 快速数据分析、ETL、嵌入式 BI、替代 Pandas 处理大文件 | 移动应用、配置存储、小型 Web 后端 | Web 后端、企业级应用、复杂事务 | 数据清洗、探索性分析(中小数据集) |
| 性能特点 | 极快的 SELECT + 聚合查询,尤其对 Parquet/CSV | 快速点查与事务 | 强一致性与复杂功能支持 | 灵活但内存消耗大,>10GB 易崩溃 |
| 语言绑定 | Python, R, C/C++, Java, JS 等 | 广泛支持 | 广泛支持 | Python 专属 |
| 是否需要服务进程 | 否(嵌入式) | 否 | 是(常驻进程) | 否 |
| 文件格式原生支持 | CSV, Parquet, JSON, Arrow | 需扩展或导入 | 需扩展(如 file_fdw) | 需加载到内存 |
1.3 DuckDB 的核心特性(嵌入式、列式存储、向量化执行)
| 特性名称 | 说明 | 注意事项 |
|---|---|---|
| 嵌入式架构 | 无须安装服务,直接通过库调用使用,启动快,部署简单。 | 不适合多用户共享访问或远程连接需求。 |
| 列式存储 | 数据按列存储,提升 I/O 效率,仅读取相关列,压缩率高。 | 插入/更新代价较高,不适合频繁写入场景。 |
| 向量化执行 | 每次处理一批数据(vector of values),减少函数调用开销,提升 CPU 利用率。 | 对现代 CPU SIMD 指令优化,性能随批大小提升。 |
| 内置文件支持 | 原生读写 Parquet、CSV、JSON 等格式,无需中间转换。 | 可直接查询外部文件,如 SELECT * FROM 'data.parquet'。 |
| SQL 兼容性 | 支持标准 SQL(包括窗口函数、CTE、JOIN 等),语法接近 PostgreSQL。 | 不支持部分高级特性如存储过程。 |
| 扩展机制 | 支持加载扩展(如 httpfs, json, parquet)以增强功能。 | 扩展需显式加载(LOAD ...)。 |
| 内存管理 | 自动管理内存,支持外部排序/聚合(spill to disk)。 | 即使数据超过内存也可处理,但速度下降。 |
1.4 安装与环境配置(Python、CLI、R 等)
| 安装方式 | 语法/命令 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| Python 安装 | pip install duckdb | 在 Python 环境中安装 DuckDB 包 | pip install duckdbpip install duckdb[parquet]pip install duckdb[icu] | 推荐添加 [parquet] 支持 Parquet 文件读写 |
| CLI 安装 | brew install duckdb (macOS) apt install duckdb (Linux) 或从官网下载二进制 | 使用命令行交互工具 | duckdbduckdb mydb.duckdb | CLI 提供交互式 SQL 环境 |
| R 安装 | install.packages(“duckdb”) | 在 R 中使用 DuckDB | install.packages("duckdb")library(duckdb) | 支持 dplyr 接口集成 |
| Node.js 安装 | npm install duckdb | 在 JavaScript/Node.js 中使用 | npm install duckdb | 功能较 Python 版有限 |
| Docker 镜像 | docker pull duckdb/duckdb | 使用容器化环境 | docker run -it duckdb/duckdb | 适合测试或 CI/CD 环境 |
| 验证安装 | import duckdb print(duckdb.version) | 检查 Python 安装是否成功 | import duckdbprint(duckdb.__version__) | 若报错请检查 Python 路径与虚拟环境 |
第2章:基础操作与连接管理
掌握如何连接、创建数据库、管理连接与执行基本命令。
2.1 连接 DuckDB(内存模式与文件模式)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| duckdb.connect() | duckdb.connect(database=':memory:' or 'path.db', read_only=False) | 创建一个到 DuckDB 数据库的连接 | conn = duckdb.connect() # 内存数据库conn = duckdb.connect('mydata.duckdb') # 文件数据库 | 默认 :memory:,关闭连接后数据丢失;指定路径则持久化 |
| 内存模式 | database=':memory:' 或不传参 | 临时数据库,速度快,重启丢失 | conn = duckdb.connect() # 等价于 ':memory:' | 适合一次性分析、测试 |
| 文件模式 | database='filename.duckdb' | 持久化数据库,数据保存在磁盘 | conn = duckdb.connect('sales.duckdb') | 文件自动创建,支持跨会话使用 |
| 只读模式 | read_only=True | 打开数据库为只读,防止误写 | conn = duckdb.connect('data.duckdb', read_only=True) | 提升安全性,适合共享数据文件 |
| 自动提交 | 默认开启 | DML 操作自动提交,无需手动 commit | conn.execute("CREATE TABLE t(x INT)") # 自动生效 | 可通过 conn.commit()/rollback() 控制事务 |
2.2 断开与关闭连接
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| close() | conn.close() | 关闭数据库连接,释放资源 | conn = duckdb.connect()conn.close() | 关闭后不可再执行语句,否则抛异常 |
| 上下文管理器 | with duckdb.connect() as conn: | 自动管理连接生命周期 | with duckdb.connect() as conn:conn.execute("CREATE TABLE t(a INT)") # 自动关闭 | 推荐做法,确保连接正确释放 |
| del conn | del conn | 删除连接对象引用 | del conn | 不保证立即关闭,仍建议显式 close() |
| 多连接管理 | 多个 conn 实例 | 并行操作不同数据库或隔离会话 | conn1 = duckdb.connect('a.duckdb')conn2 = duckdb.connect('b.duckdb') | 注意文件锁:同一文件多写会冲突 |
2.3 执行 SQL 语句的基本方法
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| execute() | conn.execute(sql) | 执行一条 SQL 语句,返回结果对象 | conn.execute("CREATE TABLE test(i INTEGER)")conn.execute("INSERT INTO test VALUES (1), (2)") | 最常用方法,支持所有 SQL 语句 |
| executescript() | conn.executescript(sql_script) | 执行多条 SQL 语句(脚本) | conn.executescript("CREATE TABLE t(x INT); INSERT INTO t VALUES (1); INSERT INTO t VALUES (2);") | 类似 SQLite,不返回中间结果 |
| df() 方法 | duckdb.sql(query).df() | 直接执行 SQL 并返回 Pandas DataFrame | result_df = duckdb.sql("SELECT * FROM test").df() | 无需先创建连接,适合快速分析 |
| arrow() 方法 | duckdb.sql(query).arrow() | 执行 SQL 返回 Arrow Table | result_arrow = duckdb.sql("SELECT * FROM test").arrow() | 零拷贝,适合与 Arrow 生态集成 |
| from_query() | duckdb.from_query(query, alias='') | 将查询封装为可复用表对象 | tbl = duckdb.from_query("SELECT i*2 AS j FROM test", 'aliased')duckdb.sql("SELECT * FROM aliased") | 便于构建复杂查询链 |
2.4 获取执行结果(fetch 与 iterate)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| fetchone() | cursor.fetchone() | 获取下一行结果,返回元组或 None | conn.execute("SELECT i FROM test")row = conn.fetchone() # (1,) | 每次调用获取一行,适合逐行处理 |
| fetchall() | cursor.fetchall() | 获取所有剩余行,返回元组列表 | rows = conn.execute("SELECT * FROM test").fetchall() # [(1,), (2,)] | 全部加载到内存,大数据慎用 |
| fetchmany(n) | cursor.fetchmany(size=n) | 获取最多 n 行,返回列表 | some_rows = conn.execute("...").fetchmany(5) | 可用于分页处理,避免内存溢出 |
| fetchdf() | cursor.fetchdf() | 获取结果为 Pandas DataFrame | df = conn.execute("SELECT * FROM test").fetchdf() | 方便与 Pandas 交互,自动推断类型 |
| fetch_arrow_table() | cursor.fetch_arrow_table() | 获取结果为 PyArrow Table | at = conn.execute("...").fetch_arrow_table() | 高效,支持零拷贝,适合大数据 |
| df()(结果对象) | result.df() | 从 sql() 结果转为 DataFrame | df = duckdb.sql("...").df() | 更简洁,推荐用于直接分析 |
| iterate() | 不适用(流式处理) | 通过 fetchone() 实现逐行迭代 | res = conn.execute("SELECT * FROM test")# while row := res.fetchone(): print(row) | 内存友好,适合超大结果集 |
提示: 推荐使用
fetchdf()或fetch_arrow_table()进行数据分析;对大型结果集,优先考虑流式fetchone()或fetchmany()。
第3章:数据定义语言(DDL)
学习如何创建、修改和删除表结构。
3.1 创建表(CREATE TABLE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| CREATE TABLE | CREATE TABLE table_name (column_def, ...) | 创建新表并定义列结构 | CREATE TABLE employees (id INTEGER, name VARCHAR, salary DOUBLE, hire_date DATE); | 必须指定列名和数据类型 |
| CREATE TABLE IF NOT EXISTS | CREATE TABLE IF NOT EXISTS ... | 避免表已存在时报错 | CREATE TABLE IF NOT EXISTS logs (ts TIMESTAMP, msg VARCHAR); | 推荐用于脚本中防止重复创建 |
| 带默认值 | column_name TYPE DEFAULT expr | 为列设置默认值 | CREATE TABLE products (id INTEGER PRIMARY KEY, name VARCHAR, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP); | DEFAULT 可使用常量或函数(如 NOW()) |
| 主键约束 | column_name TYPE PRIMARY KEY | 定义主键,确保唯一性 | CREATE TABLE users (user_id INTEGER PRIMARY KEY, email VARCHAR); | DuckDB 支持主键语法,但不强制唯一性检查(元数据用途) |
| NOT NULL 约束 | column_name TYPE NOT NULL | 约束列不允许为空 | CREATE TABLE orders (order_id INTEGER NOT NULL, amount DOUBLE NOT NULL); | 可用于数据质量控制 |
| 使用查询结果建表 | CREATE TABLE AS SELECT ... | 基于查询结果创建表 | CREATE TABLE high_earners AS SELECT * FROM employees WHERE salary > 100000; | 自动推断列名和类型,非常方便 |
| 临时表 | CREATE TEMPORARY TABLE ... | 创建会话级临时表 | CREATE TEMPORARY TABLE temp_stats AS SELECT AVG(salary) AS avg_sal FROM employees; | 仅当前连接可见,断开后自动删除 |
3.2 修改表结构(ALTER TABLE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| ADD COLUMN | ALTER TABLE table_name ADD COLUMN col_def | 添加新列 | ALTER TABLE employees ADD COLUMN department VARCHAR; | 新列对已有行的值为 NULL |
| RENAME TABLE | ALTER TABLE old_name RENAME TO new_name | 重命名表 | ALTER TABLE employees RENAME TO staff; | 表名更改,不影响数据 |
| RENAME COLUMN | ALTER TABLE table_name RENAME COLUMN old TO new | 重命名列 | ALTER TABLE employees RENAME COLUMN name TO full_name; | 需注意查询中引用的列名同步更新 |
| SET DEFAULT | ALTER TABLE table_name ALTER COLUMN col SET DEFAULT expr | 设置列默认值 | ALTER TABLE employees ALTER COLUMN department SET DEFAULT 'Unknown'; | 影响后续 INSERT 操作 |
| DROP DEFAULT | ALTER TABLE table_name ALTER COLUMN col DROP DEFAULT | 移除列默认值 | ALTER TABLE employees ALTER COLUMN department DROP DEFAULT; | 恢复为无默认值状态 |
| 支持的数据类型修改? | 不支持直接 MODIFY COLUMN | 修改列类型 | — | DuckDB 不支持直接修改列数据类型,需通过创建新表 + 数据迁移实现 |
3.3 删除表(DROP TABLE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DROP TABLE | DROP TABLE table_name | 删除表及其所有数据 | DROP TABLE employees; | 操作不可逆,请谨慎使用 |
| DROP TABLE IF EXISTS | DROP TABLE IF EXISTS table_name | 避免表不存在时报错 | DROP TABLE IF EXISTS temp_results; | 推荐用于脚本中安全删除 |
| CASCADE(暂不支持) | DROP TABLE ... CASCADE | 删除表并自动删除依赖对象 | — | DuckDB 当前不支持 CASCADE 选项 |
| RESTRICT(默认) | DROP TABLE ... RESTRICT | 仅当无依赖时才允许删除 | — | 默认行为,若存在视图等依赖会报错 |
3.4 查看表信息(PRAGMA 和 INFORMATION_SCHEMA)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| PRAGMA table_info() | PRAGMA table_info(table_name) | 查看表的列信息(名称、类型、是否主键等) | PRAGMA table_info(employees); | 类似 SQLite 的 table_info,最常用 |
| PRAGMA show() | PRAGMA show('table_name') | 显示表的全部数据(等价于 SELECT *) | PRAGMA show('employees'); | 快速预览表内容 |
| PRAGMA column_list() | PRAGMA column_list('table_name') | 列出表的所有列名 | PRAGMA column_list('employees'); | 仅返回列名 |
| 查询 INFORMATION_SCHEMA | SELECT * FROM INFORMATION_SCHEMA.COLUMNS WHERE table_name = '...' | 使用标准 SQL 获取元数据 | SELECT column_name, data_type FROM INFORMATION_SCHEMA.COLUMNS WHERE table_name = 'employees'; | 标准化方式,兼容性强 |
| duckdb_tables() | SELECT * FROM duckdb_tables() | 系统函数:列出所有表 | SELECT * FROM duckdb_tables(); | 查看当前数据库中所有表 |
| duckdb_columns() | SELECT * FROM duckdb_columns() | 系统函数:列出所有列信息 | SELECT * FROM duckdb_columns() WHERE table_name = 'employees'; | 比 PRAGMA 更灵活,可加 WHERE 条件 |
第4章:数据操作语言(DML)
掌握数据的插入、查询、更新与删除。
4.1 插入数据(INSERT INTO)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| INSERT INTO … VALUES | INSERT INTO table VALUES (v1, v2, ...) | 插入单行或多行常量值 | INSERT INTO employees VALUES (1, 'Alice', 75000, '2023-01-15');INSERT INTO employees VALUES (2, 'Bob', 80000, '2023-02-01'), (3, 'Charlie', 70000, '2023-02-10'); | 值的顺序必须与表列顺序一致 |
| INSERT INTO … (cols) VALUES | INSERT INTO table (col1, col2) VALUES (...) | 指定列插入,其余列用默认值或 NULL | INSERT INTO employees (id, name, salary) VALUES (4, 'Diana', 85000); | 更安全,避免顺序错乱 |
| INSERT INTO … SELECT | INSERT INTO table SELECT ... | 从查询结果插入数据 | INSERT INTO high_earners SELECT * FROM employees WHERE salary > 80000; | 批量迁移或聚合数据常用 |
| INSERT OR REPLACE | 不支持 | 冲突时替换 | — | DuckDB 不支持 INSERT OR REPLACE 语法 |
| 支持 UPSERT? | 不支持 ON CONFLICT | 插入或更新(UPSERT) | — | 当前版本不支持 ON CONFLICT 子句,需手动判断或使用临时表 |
4.2 查询数据(SELECT 基础)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| SELECT * FROM | SELECT * FROM table | 查询表中所有列 | SELECT * FROM employees; | 简单快速,但建议明确列名用于生产 |
| SELECT col1, col2 | SELECT col1, col2 FROM table | 查询指定列 | SELECT name, salary FROM employees; | 减少数据传输,提升性能 |
| SELECT with alias | SELECT expr AS alias FROM ... | 为列或表达式设置别名 | SELECT name AS employee_name, salary * 1.1 AS raised_salary FROM employees; | 提高可读性,用于计算字段 |
| DISTINCT | SELECT DISTINCT col FROM table | 去重查询 | SELECT DISTINCT department FROM employees; | 消除重复值,可用于单列或多列 |
| LIMIT | SELECT ... LIMIT n | 限制返回行数 | SELECT * FROM employees LIMIT 5; | 常用于调试或分页 |
| LIMIT + OFFSET | SELECT ... LIMIT n OFFSET m | 实现分页查询 | SELECT * FROM employees LIMIT 10 OFFSET 20; -- 第3页,每页10条 | OFFSET 性能随偏移量增大而下降 |
4.3 更新数据(UPDATE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| UPDATE … SET | UPDATE table SET col = value WHERE ... | 更新满足条件的行 | UPDATE employees SET salary = salary * 1.1 WHERE department = 'Engineering'; | 必须谨慎使用 WHERE,避免误更新全表 |
| 更新多列 | UPDATE table SET col1 = v1, col2 = v2 | 一次更新多个列 | UPDATE employees SET salary = salary + 5000, department = 'Senior Eng' WHERE id = 1; | 用逗号分隔多个赋值 |
| WHERE 条件 | WHERE condition | 限定更新范围 | UPDATE employees SET salary = 0 WHERE active = false; | 强烈建议始终使用 WHERE 防止全表更新 |
| 不带 WHERE | UPDATE table SET col = value | 更新所有行 | UPDATE employees SET last_check = CURRENT_DATE; | 风险高,确认后再执行 |
| 支持 RETURNING? | 不支持 RETURNING | 返回被更新的行 | — | DuckDB 当前不支持 RETURNING 子句,需额外查询验证 |
4.4 删除数据(DELETE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DELETE FROM … WHERE | DELETE FROM table WHERE condition | 删除满足条件的行 | DELETE FROM employees WHERE hire_date < '2020-01-01'; | 推荐用法,精确控制删除范围 |
| DELETE all rows | DELETE FROM table | 删除表中所有数据(保留结构) | DELETE FROM temp_cache; | 不同于 DROP TABLE,表结构仍在 |
| TRUNCATE TABLE | TRUNCATE TABLE table_name | 清空表(比 DELETE 更快) | TRUNCATE TABLE logs; | DuckDB 支持 TRUNCATE,适用于大表清空 |
| 不支持 RETURNING | DELETE ... RETURNING | 删除并返回被删行 | — | 不支持 RETURNING 子句 |
| 性能对比 | DELETE vs TRUNCATE | 清空操作性能 | — | TRUNCATE 更快,不记录逐行删除日志,推荐用于清空全表 |
第5章:数据查询进阶
深入学习 SQL 查询能力。
5.1 WHERE 条件过滤
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 比较操作符 | =, !=, <, <=, >, >= | 基本比较过滤 | SELECT * FROM employees WHERE salary >= 70000; | 支持所有标准比较运算 |
| IN 操作符 | col IN (v1, v2, ...) | 匹配值列表 | SELECT * FROM employees WHERE department IN ('Sales', 'HR'); | 等价于多个 OR 条件,性能较好 |
| NOT IN | col NOT IN (...) | 排除值列表 | SELECT * FROM employees WHERE id NOT IN (1, 2, 3); | 若列表含 NULL,结果可能为 UNKNOWN,需注意 |
| BETWEEN | col BETWEEN low AND high | 范围匹配(闭区间) | SELECT * FROM employees WHERE salary BETWEEN 50000 AND 80000; | 包含边界值,等价于 >= AND <= |
| LIKE | col LIKE 'pattern' | 模糊匹配(支持 % 和 _) | SELECT * FROM employees WHERE name LIKE 'A%'; | % 匹配任意字符,_ 匹配单个字符 |
| ILIKE | col ILIKE 'pattern' | 不区分大小写的 LIKE | SELECT * FROM employees WHERE name ILIKE 'alice'; | DuckDB 特有扩展(源自 PostgreSQL) |
| IS NULL / IS NOT NULL | col IS NULL | 判断空值 | SELECT * FROM employees WHERE department IS NULL; | 必须使用 IS NULL,不能用 = NULL |
| 逻辑操作符 | AND, OR, NOT | 组合多个条件 | SELECT * FROM employees WHERE salary > 60000 AND department = 'Engineering'; | 注意优先级,复杂表达式建议加括号 |
5.2 GROUP BY 与聚合函数
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| GROUP BY | SELECT ..., agg_func(col) FROM ... GROUP BY col | 按列分组并聚合 | SELECT department, AVG(salary) AS avg_sal FROM employees GROUP BY department; | 所有非聚合列必须出现在 GROUP BY 中 |
| 聚合函数(SUM) | SUM(expr) | 求和 | SELECT department, SUM(salary) FROM employees GROUP BY department; | 忽略 NULL 值 |
| 聚合函数(COUNT) | COUNT(*), COUNT(col) | 计数 | SELECT department, COUNT(*) AS cnt FROM employees GROUP BY department; | COUNT(*) 包含 NULL,COUNT(col) 排除 NULL |
| 聚合函数(AVG) | AVG(expr) | 平均值 | SELECT AVG(age) FROM employees; | 自动忽略 NULL |
| 聚合函数(MIN/MAX) | MIN(expr), MAX(expr) | 最小/最大值 | SELECT MIN(salary), MAX(salary) FROM employees; | 可用于数值、字符串、日期 |
| 聚合函数(STRING_AGG) | STRING_AGG(col, sep) | 字符串拼接 | SELECT department, STRING_AGG(name, ', ') FROM employees GROUP BY department; | 类似 GROUP_CONCAT |
| 多列分组 | GROUP BY col1, col2 | 按多列组合分组 | SELECT dept, role, AVG(salary) FROM employees GROUP BY dept, role; | 分组粒度更细 |
5.3 HAVING 过滤分组结果
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| HAVING | SELECT ..., agg() FROM ... GROUP BY ... HAVING condition | 对分组后的聚合结果进行过滤 | SELECT department, AVG(salary) AS avg_sal FROM employees GROUP BY department HAVING AVG(salary) > 65000; | HAVING 作用于分组后,WHERE 作用于分组前 |
| 使用聚合函数 | HAVING agg_func(col) > value | 基于聚合值过滤 | SELECT manager_id, COUNT(*) AS team_size FROM employees GROUP BY manager_id HAVING COUNT(*) >= 3; | 是 HAVING 的主要用途 |
| 多条件 HAVING | HAVING cond1 AND cond2 | 组合多个过滤条件 | HAVING AVG(salary) > 60000 AND COUNT(*) > 1; | 支持 AND/OR/NOT |
| 与 WHERE 区别 | WHERE 在 GROUP BY 前,HAVING 在后 | 不同阶段的过滤 | SELECT d, SUM(x) FROM t WHERE x > 0 GROUP BY d HAVING SUM(x) > 100; | 先用 WHERE 减少数据量,再用 HAVING 过滤结果 |
5.4 ORDER BY 与 LIMIT
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| ORDER BY | ORDER BY col [ASC|DESC] | 排序结果 | SELECT * FROM employees ORDER BY salary DESC; | 默认 ASC(升序),NULL 值排在最后 |
| 多列排序 | ORDER BY col1, col2 | 按优先级排序 | SELECT * FROM employees ORDER BY department, salary DESC; | 先按第一列排,相同再按第二列 |
| NULL 排序 | ORDER BY col NULLS FIRST|LAST | 控制 NULL 值位置 | SELECT * FROM employees ORDER BY department NULLS FIRST; | 显式控制 NULL 排序行为 |
| LIMIT | LIMIT n | 限制返回行数 | SELECT * FROM employees ORDER BY salary DESC LIMIT 5; | 常用于 Top-N 查询 |
| LIMIT + OFFSET | LIMIT n OFFSET m | 实现分页 | SELECT * FROM employees ORDER BY id LIMIT 10 OFFSET 20; | OFFSET 性能随偏移增大而下降 |
| FETCH FIRST | FETCH FIRST n ROWS ONLY | 标准化 LIMIT 替代语法 | SELECT * FROM employees ORDER BY salary DESC FETCH FIRST 5 ROWS ONLY; | SQL 标准语法,与 LIMIT 等价 |
5.5 JOIN 多表连接(INNER, LEFT, RIGHT, FULL)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| INNER JOIN | SELECT ... FROM A INNER JOIN B ON A.key = B.key | 内连接:仅保留匹配行 | SELECT e.name, d.dept_name FROM employees e INNER JOIN departments d ON e.dept_id = d.id; | 最常用,只返回两表都有的匹配记录 |
| LEFT JOIN | SELECT ... FROM A LEFT JOIN B ON ... | 左外连接:保留左表所有行 | SELECT e.name, COALESCE(d.dept_name, 'Unassigned') FROM employees e LEFT JOIN departments d ON e.dept_id = d.id; | 左表全保留,右表无匹配则为 NULL |
| RIGHT JOIN | SELECT ... FROM A RIGHT JOIN B ON ... | 右外连接:保留右表所有行 | SELECT e.name, d.dept_name FROM employees e RIGHT JOIN departments d ON e.dept_id = d.id; | 右表全保留,左表无匹配则为 NULL |
| FULL JOIN | SELECT ... FROM A FULL JOIN B ON ... | 全外连接:保留两表所有行 | SELECT e.name, d.dept_name FROM employees e FULL JOIN departments d ON e.dept_id = d.id; | 任一表有记录即保留,缺失补 NULL |
| USING 子句 | JOIN ... USING (col) | 简化等值连接(列名相同) | SELECT * FROM emp USING (dept_id); | 当连接列名相同时可简化语法 |
| 自连接 | JOIN 同一表 | 表自身连接 | SELECT a.name, b.name AS manager FROM employees a JOIN employees b ON a.manager_id = b.id; | 用于层级关系(如员工-经理) |
| 多表 JOIN | 链式 JOIN | 连接三个及以上表 | SELECT e.name, d.name, p.project_name FROM emp e JOIN dept d ON e.d = d.id JOIN projects p ON e.p_id = p.id; | 按业务逻辑顺序连接 |
5.6 子查询与 CTE(WITH 子句)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 标量子查询 | SELECT (SELECT ...) FROM ... | 返回单值的子查询 | SELECT name, (SELECT AVG(salary) FROM employees) AS avg_all FROM employees; | 必须返回单行单列,常用于计算字段 |
| WHERE 中子查询 | WHERE col OP (SELECT ...) | 在 WHERE 中使用子查询 | SELECT * FROM employees WHERE salary > (SELECT AVG(salary) FROM employees); | 支持 IN, EXISTS, 比较等 |
| FROM 中子查询 | SELECT ... FROM (SELECT ...) AS alias | 衍生表(内联视图) | SELECT dept, avg_sal FROM (SELECT department AS dept, AVG(salary) AS avg_sal FROM employees GROUP BY department) t WHERE avg_sal > 60000; | 子查询必须有别名 |
| WITH 子句(CTE) | WITH cte_name AS (query) SELECT ... FROM cte_name | 公共表表达式,提高可读性 | WITH dept_avg AS (SELECT department, AVG(salary) AS avg_sal FROM employees GROUP BY department) SELECT * FROM dept_avg WHERE avg_sal > 65000; | 可定义多个 CTE,支持递归(见下) |
| 递归 CTE | WITH RECURSIVE cte AS (base_query UNION ALL recursive_query) | 实现递归查询(如树形结构) | WITH RECURSIVE emp_tree AS (SELECT id, name, manager_id, 1 AS level FROM employees WHERE manager_id IS NULL UNION ALL SELECT e.id, e.name, e.manager_id, et.level + 1 FROM employees e JOIN emp_tree et ON e.manager_id = et.id) SELECT * FROM emp_tree; | DuckDB 支持递归 CTE,用于组织架构、路径遍历等 |
5.7 窗口函数(OVER, PARTITION BY, ROWS/RANGE)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| ROW_NUMBER() | ROW_NUMBER() OVER (ORDER BY col) | 行号(唯一) | SELECT name, salary, ROW_NUMBER() OVER (ORDER BY salary DESC) AS rank FROM employees; | 每行分配唯一序号,即使值相同也不同 |
| RANK() | RANK() OVER (ORDER BY col) | 排名(跳跃) | RANK() OVER (ORDER BY salary DESC) AS ranked | 相同值同名次,下一名跳过(如 1,1,3) |
| DENSE_RANK() | DENSE_RANK() OVER (...) | 密集排名 | DENSE_RANK() OVER (ORDER BY salary DESC) AS dense_rank | 相同值同名次,下一名连续(如 1,1,2) |
| PARTITION BY | OVER (PARTITION BY col ORDER BY col2) | 分组窗口计算 | SELECT name, department, salary, AVG(salary) OVER (PARTITION BY department) AS dept_avg FROM employees; | 在每个分区内独立计算 |
| LAG() / LEAD() | LAG(col, offset, default) | 访问前/后 N 行 | SELECT date, revenue, LAG(revenue, 1) OVER (ORDER BY date) AS prev_rev FROM sales; | 时间序列分析常用 |
| SUM() OVER | SUM(col) OVER (...) | 累计求和 | SUM(sales) OVER (ORDER BY date ROWS UNBOUNDED PRECEDING) AS cum_sales | 实现运行总计 |
| ROWS / RANGE | OVER (ORDER BY ... ROWS BETWEEN ...) | 定义窗口范围 | SUM(sales) OVER (ORDER BY date ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) | ROWS 按物理行,RANGE 按值范围 |
| FIRST_VALUE() / LAST_VALUE() | FIRST_VALUE(col) OVER (...) | 获取窗口内首/尾值 | FIRST_VALUE(salary) OVER (PARTITION BY dept ORDER BY hire_date) | 常用于获取最早/最晚记录 |
第6章:数据类型与函数
熟悉 DuckDB 支持的数据类型与内置函数。
6.1 基本数据类型(INTEGER, VARCHAR, DATE, BOOLEAN 等)
| 数据类型 | 说明 | 注意事项 |
|---|---|---|
| BOOLEAN | 布尔值:TRUE / FALSE / NULL | 支持标准逻辑运算 |
| TINYINT | 8位整数,范围 -128 到 127 | 小整数存储,节省空间 |
| SMALLINT | 16位整数,范围 -32,768 到 32,767 | — |
| INTEGER | 32位整数,范围约 ±21亿 | 最常用整型 |
| BIGINT | 64位整数 | 大数、时间戳等 |
| HUGEINT | 128位整数(实验性) | 极大数值,如金融计算 |
| REAL / FLOAT | 32位浮点数 | 单精度 |
| DOUBLE / FLOAT8 | 64位浮点数 | 双精度,推荐用于科学计算 |
| DECIMAL(p,s) | 精确小数,p=精度,s=标度 | DECIMAL(10,2) 表示最多10位,2位小数,适合货币 |
| VARCHAR / TEXT | 可变长度字符串 | 无默认长度限制 |
| CHAR(n) | 固定长度字符串 | 不足补空格 |
| DATE | 日期(年-月-日) | 范围 '0001-01-01' 到 '9999-12-31' |
| TIME | 时间(时:分:秒[.毫秒]) | 精确到微秒 |
| TIMESTAMP | 日期+时间 | 默认无时区 |
| TIMESTAMP WITH TIME ZONE | 带时区的时间戳 | 存储 UTC,显示时转换 |
| INTERVAL | 时间间隔 | 如 INTERVAL '1 day',用于日期运算 |
6.2 复合类型(STRUCT, LIST, MAP)
| 复合类型 | 语法/构造方式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| STRUCT | ROW(expr, ...) 或 {'key': val, ...} | 类似记录或对象 | SELECT ROW('Alice', 30) AS person;SELECT {'name': 'Bob', 'age': 25} AS info; | 可嵌套,通过 col.key 访问字段 |
| LIST / ARRAY | [elem1, elem2, ...] | 有序集合 | SELECT [1, 2, 3] AS nums;SELECT ['a', 'b'] AS tags; | 支持数组操作函数 |
| MAP | MAP{'key': val, ...} | 键值对集合 | SELECT MAP{'x': 1, 'y': 2} AS coords; | 类似字典,键必须唯一 |
| 访问 STRUCT 字段 | struct_col.key | 提取结构体字段 | SELECT info.name FROM (SELECT {'name':'Alice'} AS info) t; | 点号语法 |
| 访问 LIST 元素 | list_col[index] | 按索引访问(从1开始) | SELECT tags[1] FROM mytable; -- 第一个元素 | 索引越界返回 NULL |
| LIST 聚合 | LIST_VALUE(...) 或 LIST(...) | 聚合生成列表 | SELECT department, LIST(name) AS members FROM employees GROUP BY department; | 强大功能,用于分组聚合 |
| MAP 聚合 | MAP(keys, values) | 从两列创建 MAP | SELECT MAP(ARRAY['a','b'], ARRAY[1,2]) AS m; -- {'a':1, 'b':2} | 键值数组长度需一致 |
6.3 日期时间函数(DATE_TRUNC, INTERVAL, EXTRACT 等)
| 函数名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| CURRENT_DATE | CURRENT_DATE | 返回当前日期 | SELECT CURRENT_DATE; | 格式:YYYY-MM-DD |
| CURRENT_TIME | CURRENT_TIME | 返回当前时间(含时区) | SELECT CURRENT_TIME; | 包含时区偏移 |
| CURRENT_TIMESTAMP | CURRENT_TIMESTAMP | 返回当前时间戳 | SELECT CURRENT_TIMESTAMP; | 最常用,包含日期和时间 |
| NOW() | NOW() | 同 CURRENT_TIMESTAMP | SELECT NOW(); | PostgreSQL 兼容别名 |
| DATE_TRUNC | DATE_TRUNC(unit, timestamp) | 截断时间到指定粒度 | SELECT DATE_TRUNC('month', '2023-04-15 12:34:56'); -- 结果: 2023-04-01 00:00:00 | 支持 'year', 'month', 'day', 'hour', 'minute', 'second' |
| EXTRACT | EXTRACT(field FROM timestamp) | 提取时间部分 | SELECT EXTRACT(YEAR FROM '2023-04-15') AS y;SELECT EXTRACT(HOUR FROM CURRENT_TIMESTAMP); | field 可为 YEAR, MONTH, DAY, HOUR, MINUTE, SECOND, DOW, DOY 等 |
| DATE_PART | DATE_PART(unit, timestamp) | 类似 EXTRACT,但 unit 为字符串 | SELECT DATE_PART('day', '2023-04-15'); | 与 EXTRACT 功能重叠,PostgreSQL 风格 |
| INTERVAL | INTERVAL 'value unit' | 创建时间间隔 | SELECT '2023-01-01'::DATE + INTERVAL '1 day'; -- 结果: 2023-01-02 | 支持 'day', 'month', 'year', 'hour', 'minute' 等 |
| 年月日构造 | MAKE_DATE(year, month, day) | 从年月日构造日期 | SELECT MAKE_DATE(2023, 4, 15); | 安全构造日期,自动校验 |
| 时间加减 | datetime ± INTERVAL | 日期时间运算 | SELECT CURRENT_TIMESTAMP - INTERVAL '7 days'; | 支持加减各种单位 |
| AGE | AGE(timestamp) | 计算与当前时间的间隔 | SELECT AGE('1990-05-20'); | 返回 INTERVAL 类型,表示年龄 |
6.4 字符串处理函数(SUBSTR, REPLACE, REGEXP 等)
| 函数名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| LENGTH / CHAR_LENGTH | LENGTH(string) | 返回字符串字符数 | SELECT LENGTH('Hello'); -- 5 | 支持多字节字符(如中文) |
| LOWER / UPPER | LOWER(string), UPPER(string) | 大小写转换 | SELECT UPPER('hello'); -- 'HELLO' | 基本文本处理 |
| TRIM | TRIM(string) | 去除首尾空格 | SELECT TRIM(' abc '); -- 'abc' | 支持 LTRIM, RTRIM 去除单侧 |
| SUBSTR / SUBSTRING | SUBSTR(string, start, length?) | 截取子串 | SELECT SUBSTR('Hello', 2, 3); -- 'ell'SELECT SUBSTR('Hello', -3); -- 'llo' | 起始位置从1开始,支持负数(从末尾) |
| LEFT / RIGHT | LEFT(string, n), RIGHT(string, n) | 取左/右 N 个字符 | SELECT LEFT('Hello', 2); -- 'He'SELECT RIGHT('Hello', 3); -- 'llo' | 简化常用操作 |
| REPLACE | REPLACE(string, old, new) | 字符串替换 | SELECT REPLACE('Hello World', 'World', 'DuckDB'); -- 'Hello DuckDB' | 全局替换所有匹配 |
| REVERSE | REVERSE(string) | 字符串反转 | SELECT REVERSE('abc'); -- 'cba' | 文本处理工具 |
| CONCAT | string1 || string2 或 CONCAT(s1,s2,...) | 字符串拼接 | SELECT 'Hello ' || 'World'; | || 是标准 SQL 拼接操作符 |
| SPLIT_PART | SPLIT_PART(string, delim, n) | 按分隔符分割取第 n 部分 | SELECT SPLIT_PART('a,b,c', ',', 2); -- 'b' | 类似 Python 的 split()[n-1] |
| REGEXP_MATCHES | REGEXP_MATCHES(string, pattern) | 正则匹配,返回匹配组 | SELECT REGEXP_MATCHES('abc123', '([a-z]+)([0-9]+)'); -- ['abc', '123'] | 返回 LIST,支持捕获组 |
| REGEXP_EXTRACT | REGEXP_EXTRACT(string, pattern, group?) | 提取正则匹配的指定组 | SELECT REGEXP_EXTRACT('User: Alice', 'User: (\w+)', 1); -- 'Alice' | group 缺省为1 |
| REGEXP_REPLACE | REGEXP_REPLACE(string, pattern, replacement) | 正则替换 | SELECT REGEXP_REPLACE('abc123', '\d+', 'XXX'); -- 'abcXXX' | 支持全局替换 |
| LIKE / ILIKE | string LIKE 'pattern' | 模式匹配(见 5.1) | SELECT 'hello' LIKE 'h%'; -- true | % 任意,_ 单字符 |
6.5 数值与聚合函数(SUM, AVG, ROUND 等)
| 函数名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| SUM | SUM(expr) | 求和 | SELECT SUM(salary) FROM employees; | 忽略 NULL,返回对应类型 |
| AVG | AVG(expr) | 平均值 | SELECT AVG(price) FROM products; | 自动转换为 DOUBLE |
| COUNT | COUNT(*), COUNT(expr) | 计数 | SELECT COUNT(*) FROM employees;SELECT COUNT(commission) FROM employees; -- 非 NULL 个数 | COUNT(*) 统计行数 |
| MIN / MAX | MIN(expr), MAX(expr) | 最小/最大值 | SELECT MIN(age), MAX(age) FROM people; | 支持数值、字符串、日期 |
| ROUND | ROUND(expr, decimals?) | 四舍五入 | SELECT ROUND(3.14159, 2); -- 3.14SELECT ROUND(123.456, -1); -- 120.0 | decimals 可为负数(小数点左) |
| CEIL / FLOOR | CEIL(x), FLOOR(x) | 向上/下取整 | SELECT CEIL(3.2), FLOOR(3.8); -- 4, 3 | 返回 DOUBLE |
| ABS | ABS(x) | 绝对值 | SELECT ABS(-10); -- 10 | — |
| MOD | x % y 或 MOD(x, y) | 取模运算 | SELECT 10 % 3; -- 1 | — |
| POWER | POWER(x, y) | 幂运算 | SELECT POWER(2, 3); -- 8 | — |
| SQRT | SQRT(x) | 平方根 | SELECT SQRT(16); -- 4 | — |
| LOG / LN | LOG(x), LN(x) | 对数(10为底 / 自然) | SELECT LOG(100); -- 2.0SELECT LN(2.718); -- ~1.0 | — |
| GREATEST / LEAST | GREATEST(a,b,...), LEAST(a,b,...) | 返回最大/最小值 | SELECT GREATEST(3, 7, 1); -- 7 | 可比较多个值 |
| COVAR_SAMP / CORR | COVAR_SAMP(x,y), CORR(x,y) | 样本协方差 / 相关系数 | SELECT CORR(income, spending) FROM data; | 统计分析函数 |
| STDDEV / VAR | STDDEV_SAMP(expr), VAR_SAMP(expr) | 样本标准差 / 方差 | SELECT STDDEV_SAMP(score) FROM exam_results; | 支持 _POP 版本(总体) |
6.6 条件函数(CASE, COALESCE, IF 等)
| 函数名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| CASE(简单) | CASE expr WHEN val THEN result ... END | 简单分支 | SELECT name, CASE department WHEN 'Sales' THEN 'A' WHEN 'Eng' THEN 'B' ELSE 'C' END AS grade FROM employees; | 类似 switch |
| CASE(搜索) | CASE WHEN cond THEN result ... END | 条件分支 | SELECT name, salary, CASE WHEN salary < 50000 THEN 'Low' WHEN salary < 80000 THEN 'Medium' ELSE 'High' END AS level FROM employees; | 更灵活,支持任意条件 |
| COALESCE | COALESCE(expr1, expr2, ...) | 返回第一个非 NULL 值 | SELECT COALESCE(NULL, 'default', 'backup'); -- 'default'SELECT COALESCE(phone, mobile, 'N/A') FROM contacts; | 常用于空值替换 |
| IF | IF(condition, true_val, false_val) | 三元条件判断 | SELECT IF(score >= 60, 'Pass', 'Fail') FROM exams; | 简化简单条件逻辑 |
| NULLIF | NULLIF(expr1, expr2) | 若两值相等则返回 NULL | SELECT NULLIF(0, 0); -- NULLSELECT NULLIF(1, 0); -- 1 | 常用于避免除零:x / NULLIF(y, 0) |
| ISNULL / IFNULL | ISNULL(expr, replacement) | 扩展:若为 NULL 则替换 | SELECT ISNULL(department, 'Unknown') FROM employees; | SQLite 兼容函数,等价于 COALESCE |
| AND / OR / NOT | condition AND condition | 逻辑运算 | WHERE active AND (age > 18 OR override); | 标准布尔逻辑 |
| BETWEEN / IN | expr BETWEEN low AND high | 范围/集合判断 | WHERE score BETWEEN 0 AND 100;WHERE status IN ('active', 'pending'); | 见 5.1 |
第7章:外部数据交互
学习如何导入导出数据,连接外部文件与系统。
7.1 读取 CSV 文件(READ_CSV)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| READ_CSV() | SELECT * FROM READ_CSV('file.csv') | 直接查询 CSV 文件 | SELECT * FROM READ_CSV('employees.csv'); | DuckDB 自动推断模式 |
| 指定列类型 | READ_CSV('file.csv', columns={'col': 'type', ...}) | 显式定义列和类型 | SELECT * FROM READ_CSV('data.csv', columns={'id': 'INTEGER', 'name': 'VARCHAR', 'salary': 'DOUBLE'}); | 避免类型推断错误 |
| 自动模式推断 | AUTO_DETECT=TRUE | 启用自动检测分隔符、标题等 | SELECT * FROM READ_CSV('data.txt', AUTO_DETECT=TRUE); | 推荐用于格式不规范的文件 |
| 指定分隔符 | DELIMITER=',' | 设置字段分隔符 | SELECT * FROM READ_CSV('data.tsv', DELIMITER='\t'); | 支持 ',' 等 |
| 是否有标题行 | HEADER=TRUE|FALSE | 指定文件是否包含列名 | SELECT * FROM READ_CSV('no_header.csv', HEADER=FALSE); | — |
| 多文件读取 | READ_CSV(['f1.csv','f2.csv']) | 合并多个 CSV 文件 | SELECT * FROM READ_CSV(['sales_2023.csv', 'sales_2024.csv']); | 自动 UNION ALL,结构需一致 |
| 通配符读取 | READ_CSV('sales_*.csv') | 使用通配符匹配文件 | SELECT region, SUM(amount) FROM READ_CSV('sales_*.csv') GROUP BY region; | 方便批量处理 |
| 创建表基于 CSV | CREATE TABLE t AS SELECT * FROM READ_CSV(...) | 将 CSV 数据加载到表中 | CREATE TABLE employees AS SELECT * FROM READ_CSV('emp.csv', AUTO_DETECT=TRUE); | 实现”导入”操作 |
7.2 读取 Parquet 文件(READ_PARQUET)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| READ_PARQUET() | SELECT * FROM READ_PARQUET('file.parquet') | 查询 Parquet 文件 | SELECT * FROM READ_PARQUET('users.parquet'); | 高效列式存储,推荐使用 |
| 读取目录 | READ_PARQUET('dir/') | 读取目录下所有 Parquet 文件 | SELECT * FROM READ_PARQUET('logs/'); | 自动合并,支持分区目录 |
| 通配符读取 | READ_PARQUET('data_*.parquet') | 匹配多个文件 | SELECT * FROM READ_PARQUET('sales_*.parquet'); | 类似 CSV |
| 元数据查看 | PRAGMA file_metadata('file.parquet') | 查看 Parquet 文件元数据 | PRAGMA file_metadata('users.parquet'); | 包含行组、列统计等信息 |
| 列裁剪 | SELECT col1, col2 FROM ... | 自动只读取所需列 | SELECT name, age FROM READ_PARQUET('big.parquet'); | Parquet 优势:I/O 高效 |
| 谓词下推 | WHERE col = value | 自动下推过滤条件 | SELECT * FROM READ_PARQUET('data.parquet') WHERE status = 'active'; | 仅扫描匹配的行组,性能极佳 |
| 分区读取 | 支持 Hive 分区 | 自动识别分区列 | -- 目录结构: logs/year=2023/month=01/file.parquetSELECT * FROM READ_PARQUET('logs/') WHERE year = 2023; | 分区列自动作为字段可用 |
| 创建表基于 Parquet | CREATE TABLE t AS SELECT * FROM READ_PARQUET(...) | 导入 Parquet 到表 | CREATE TABLE users AS SELECT * FROM READ_PARQUET('users.parquet'); | 快速加载 |
7.3 读取 JSON 文件(READ_JSON)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| READ_JSON() | SELECT * FROM READ_JSON('file.json') | 读取 JSON 文件(数组) | SELECT * FROM READ_JSON('data.json'); | 支持 JSON Lines (.jsonl) 和数组 |
| JSON Lines | 每行一个 JSON 对象 | 逐行解析 | SELECT * FROM READ_JSON('logs.jsonl'); | 流式处理,内存友好 |
| 自动模式推断 | AUTO_DETECT=TRUE | 推断嵌套结构类型 | SELECT * FROM READ_JSON('data.json', AUTO_DETECT=TRUE); | 处理复杂 JSON 结构 |
| 指定模式 | COLUMNS={'col': 'type', ...} | 显式定义模式 | SELECT * FROM READ_JSON('data.json', COLUMNS={'id': 'INT', 'info': 'STRUCT(name VARCHAR, age INT)'}); | 精确控制,避免推断错误 |
| 嵌套 JSON 访问 | 使用 -> 操作符 | 提取嵌套字段 | SELECT json_col->'address'->>'city' AS city FROM (SELECT READ_JSON('users.json') AS json_col); | -> 返回 JSON,->> 返回文本 |
| 读取多个文件 | READ_JSON(['f1.json','f2.json']) | 合并多个 JSON 文件 | SELECT * FROM READ_JSON(['data1.json', 'data2.json']); | 结构应一致 |
| 创建表基于 JSON | CREATE TABLE t AS SELECT * FROM READ_JSON(...) | 导入 JSON 数据 | CREATE TABLE events AS SELECT * FROM READ_JSON('events.jsonl', AUTO_DETECT=TRUE); | 实现 JSON 到关系表转换 |
7.4 导出为 CSV、Parquet、JSON
| 格式 | 导出语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 导出为 CSV | COPY (query) TO 'file.csv' WITH (FORMAT CSV, HEADER true) | 将查询结果保存为 CSV | COPY (SELECT * FROM employees) TO 'output.csv' WITH (FORMAT CSV, HEADER true); | 支持 HEADER, DELIMITER, QUOTE 等选项 |
| 导出为 Parquet | COPY (query) TO 'file.parquet' (FORMAT PARQUET) | 导出为 Parquet 格式 | COPY (SELECT * FROM sales) TO 'sales.parquet'; | 高效存储,推荐长期保存 |
| 导出为 JSON | COPY (query) TO 'file.json' (FORMAT JSON) | 导出为 JSON 数组 | COPY (SELECT * FROM users) TO 'users.json'; | 生成 [{},{}] 格式 |
| 导出为 JSON Lines | COPY ... (FORMAT JSON, ARRAY false) | 每行一个 JSON 对象 | COPY (SELECT * FROM logs) TO 'logs.jsonl' WITH (FORMAT JSON, ARRAY false); | 流式处理友好 |
| 压缩导出 | COMPRESSION 'zstd' 等 | 启用压缩 | COPY (SELECT * FROM big_table) TO 'data.parquet' (COMPRESSION 'ZSTD'); | 支持 SNAPPY, GZIP, ZSTD 等,减小文件大小 |
| 分区导出 | PARTITION_BY (col) | 按列分区导出 | COPY (SELECT * FROM sales) TO 'sales/' (FORMAT PARQUET) PARTITION_BY (region); | 生成 Hive 分区目录结构,便于后续查询 |
| 从表导出 | COPY table_name TO ... | 直接导出整个表 | COPY employees TO 'emp.csv' (FORMAT CSV, HEADER true); | 简化操作 |
7.5 与 Pandas DataFrame 交互(df.to_sql / duckdb.sql)
| 方法/方向 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DuckDB 查询 → DataFrame | conn.execute(sql).df() | 执行 SQL 并返回 DataFrame | import duckdbcon = duckdb.connect()df = con.execute("SELECT * FROM employees").df() | 最常用方式 |
| DataFrame → DuckDB 表 | con.register('name', df) | 注册 DataFrame 为虚拟表 | con.register('sales_df', sales_dataframe)result = con.execute("SELECT region, SUM(amount) FROM sales_df GROUP BY region").df() | 无需复制数据,高效 |
| 直接查询 DataFrame | con.execute("SELECT * FROM df").df() | 把 DataFrame 当作表查询 | # 假设 df 已存在filtered = con.execute("SELECT * FROM df WHERE age > 30").df() | DataFrame 名称作为表名 |
| 将结果保存到 DataFrame | assign 或 create_view | 创建视图或变量 | con.execute("CREATE VIEW v AS SELECT * FROM df WHERE active")v_df = con.table('v').df() | 用于复杂工作流 |
| 批量处理 | 结合 register 和 SQL | 利用 DuckDB 处理大 DataFrame | con.register('large_df', big_df)aggregated = con.execute("SELECT category, AVG(value) AS avg_val FROM large_df GROUP BY category").df() | 避免 Pandas 内存瓶颈 |
7.6 与 Arrow 表集成
| 方法/方向 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| Arrow Table → DuckDB | con.register('name', table) | 注册 Arrow 表为虚拟表 | import pyarrow as paimport duckdbtable = pa.table([pa.array([1, 2, 3]), pa.array(['a', 'b', 'c'])], names=['id', 'value'])con = duckdb.connect()con.register('arrow_table', table)result = con.execute("SELECT * FROM arrow_table WHERE id > 1").fetch_arrow_table() | 零拷贝,极高性能 |
| DuckDB → Arrow Table | conn.execute(sql).fetch_arrow_table() | 执行 SQL 返回 Arrow 表 | arrow_result = con.execute("SELECT * FROM employees").fetch_arrow_table() | 用于与 Arrow 生态(如 Polars)集成 |
| 批量数据交换 | register + fetch_arrow_table | 在 DuckDB 和 Arrow 间高效传输 | # 处理 Arrow 数据processed = con.execute("SELECT id, UPPER(value) AS upper_val FROM arrow_table").fetch_arrow_table() | 避免序列化开销 |
| 支持数据类型 | 完全兼容 Arrow 类型系统 | 无缝映射 | — | 包括 STRUCT, LIST, DICTIONARY 等复杂类型 |
| 内存效率 | 零拷贝访问 | 不复制数据 | — | 特别适合大内存数据集处理 |
| 与 Pandas 对比 | 更高效 | Arrow 是列式内存格式 | — | 对于数值计算和大表,Arrow 比 Pandas DataFrame 更高效 |
第8章:性能优化与高级特性
提升查询效率,掌握 DuckDB 高级功能。
8.1 索引支持(CREATE INDEX)
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 创建 B-Tree 索引 | CREATE INDEX idx_name ON table(col) | 加速等值和范围查询 | CREATE INDEX idx_salary ON employees(salary); | DuckDB 支持 B-Tree 索引,适用于 =, >, <, BETWEEN, IN |
| 创建复合索引 | CREATE INDEX idx_multi ON table(col1, col2) | 支持多列查询条件 | CREATE INDEX idx_dept_role ON employees(department, role); | 遵循最左前缀原则,查询需使用索引前列 |
| 唯一索引 | CREATE UNIQUE INDEX idx_uniq ON table(col) | 确保列值唯一并加速查询 | CREATE UNIQUE INDEX idx_email ON employees(email); | 插入重复值会报错 |
| 查看索引 | PRAGMA show_indexes ON table_name | 列出表的所有索引 | PRAGMA show_indexes ON employees; | 用于调试和验证 |
| 删除索引 | DROP INDEX idx_name | 移除索引 | DROP INDEX idx_salary; | 重建或优化时使用 |
| 索引与性能 | 在 WHERE、JOIN、ORDER BY 中使用 | 减少扫描行数,加速排序 | SELECT * FROM employees WHERE salary > 70000; -- 使用 idx_salary | 索引会增加写入开销,需权衡 |
| 自动索引(临时) | DuckDB 可能为优化 JOIN 创建临时索引 | 内部优化机制 | — | 用户无需干预,但了解有助于理解执行计划 |
8.2 分区表与分块读取
| 方法/语法 | 语法格式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 创建分区表 | CREATE TABLE ... (col TYPE) PARTITION BY (partition_col) | 按列对表进行物理分区 | CREATE TABLE sales (sale_date DATE, amount DOUBLE, region VARCHAR) PARTITION BY (region); | 数据按分区列的值存储在不同文件/块中 |
| 插入分区数据 | INSERT INTO table VALUES (...) | 自动路由到对应分区 | INSERT INTO sales VALUES ('2023-01-01', 1000, 'North'); | 写入时自动定位分区 |
| 分区裁剪 | WHERE partition_col = value | 查询时自动跳过无关分区 | SELECT SUM(amount) FROM sales WHERE region = 'South'; | 极大减少 I/O,核心性能优势 |
| 外部文件分区 | READ_PARQUET('path/', hive_partitioning=TRUE) | 读取 Hive 风格分区目录 | SELECT * FROM READ_PARQUET('sales_data/', hive_partitioning=TRUE) WHERE year = 2023 AND month = '01'; | 目录结构如 sales_data/year=2023/month=01/file.parquet |
| 分块读取(内部) | DuckDB 自动分块 | 将大文件/表分块处理 | — | 内部机制,支持并行处理和内存管理 |
| 优势 | 减少扫描范围、并行处理、高效存储 | 优化大表查询性能 | — | 特别适合时间序列、日志、按区域/类别划分的数据 |
| 注意事项 | 分区列选择 | 选择高基数、常用于过滤的列 | — | 避免过度分区(小文件问题) |
8.3 并行查询与向量化执行机制
| 特性 | 说明 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 并行查询 | DuckDB 自动利用多核 CPU | 加速大数据集处理 | SELECT department, AVG(salary), COUNT(*) FROM employees GROUP BY department; | 大多数查询自动并行化,无需手动配置 |
| 向量化执行 | 一次处理一批数据(向量) | 提高 CPU 缓存效率和指令吞吐 | — | 核心引擎特性,所有操作均向量化 |
| 批处理大小 | 默认 1024 行/批 | 平衡内存和计算效率 | — | 可通过 PRAGMA vector_size 调整(不推荐) |
| 并行 I/O | 读取 Parquet/CSV 时并行 | 快速加载外部数据 | SELECT COUNT(*) FROM READ_PARQUET('huge_file.parquet'); | 利用多核和磁盘带宽 |
| 并行聚合 | GROUP BY 在多个线程上并行执行 | 加速分组聚合 | — | 对于高基数分组效果显著 |
| 并行排序 | ORDER BY 多线程执行 | 加速大型排序 | SELECT * FROM big_table ORDER BY key; | 使用外部排序算法,支持大于内存的数据 |
| 资源控制 | 可设置线程数 | 限制资源使用 | PRAGMA threads=4; -- 设置最大线程数 | 在资源受限环境中有用 |
8.4 使用 EXPLAIN 查看执行计划
| 命令 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| EXPLAIN | EXPLAIN SELECT ... | 查看查询执行计划(文本) | EXPLAIN SELECT * FROM employees WHERE salary > 70000; | 显示操作符树,了解执行流程 |
| EXPLAIN QUERY PLAN | EXPLAIN QUERY PLAN SELECT ... | 更简洁的执行计划 | EXPLAIN QUERY PLAN SELECT COUNT(*) FROM sales; | 类似 SQLite,显示高层次操作 |
| 输出解读 | 包含操作符如 SEQ_SCAN, HASH_JOIN, AGGREGATE | 识别性能瓶颈 | — | 关注是否使用索引、JOIN 策略、聚合方式 |
| 识别全表扫描 | SEQ_SCAN 表名 | 可能缺少索引或谓词下推失败 | — | 检查 WHERE 条件和索引 |
| 识别 JOIN 类型 | NESTED LOOP, HASH JOIN, MERGE JOIN | 了解连接效率 | — | HASH JOIN 通常最快,MERGE JOIN 用于有序数据 |
| 结合性能测试 | 配合实际执行时间 | 验证优化效果 | PRAGMA enable_profiling;SELECT ...; -- 查看详细耗时 | enable_profiling 可生成更详细的性能报告 |
8.5 缓存与内存管理
| 特性/方法 | 说明 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 内存数据库 | 默认在内存中运行 | 高速数据处理 | con = duckdb.connect() # 内存数据库 | 断开连接后数据丢失 |
| 持久化数据库 | connect('db_name.db') | 将数据持久化到磁盘 | con = duckdb.connect('my_db.db') | 数据库文件可复用 |
| 内存限制 | PRAGMA memory_limit | 设置最大内存使用 | PRAGMA memory_limit='2GB'; | 防止 OOM,处理超大数据集时启用外部排序/聚合 |
| 自动溢出到磁盘 | 当内存不足时 | 支持处理大于内存的数据集 | — | 透明进行,性能会下降但保证查询完成 |
| 数据缓存 | DuckDB 自动缓存热数据 | 提高重复查询性能 | — | 无需手动干预 |
| 结果缓存 | 无内置查询缓存 | 每次执行重新计算 | — | 应用层可实现缓存逻辑 |
| 优化建议 | 合理设置 memory_limit,利用持久化 | 平衡性能与资源 | — | 对于 ETL 任务,内存数据库 + 最后持久化是高效模式 |
8.6 自定义函数(UDF)与扩展(如 httpfs, json, parquet)
| 特性/方法 | 语法/方式 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 加载扩展 | LOAD 'extension_name' | 启用额外功能 | LOAD 'httpfs'; -- 访问 S3/GCSLOAD 'json'; -- 增强 JSON 支持LOAD 'parquet'; -- 确保 Parquet 支持 | 大多数功能默认启用,httpfs 需显式加载 |
| HTTPFS 扩展 | 读取云存储文件 | 直接查询 S3、GCS 等 | SELECT * FROM READ_PARQUET('s3://bucket/data/*.parquet'); | 需配置凭据(环境变量或 AWS CLI) |
| 创建标量 UDF(Python) | con.create_function() | 在 Python 中定义自定义函数 | def greet(name): return f"Hello, {name}!"con.create_function('greet', greet, [str], str)con.execute("SELECT greet('Alice')").fetchone() # 'Hello, Alice!' | 函数在 DuckDB 内部调用,可用于 SQL |
| 创建向量化 UDF | 使用 Arrow 集成 | 高性能自定义函数 | import pyarrow.compute as pcdef add_ten(arr): return pc.add(arr, 10)con.create_vectorized_function('add_ten', add_ten, ['DOUBLE'], 'DOUBLE') | 利用 Arrow 计算,性能接近原生 |
| 注册聚合 UDF | create_aggregate_function() | 定义自定义聚合函数 | — | 较复杂,需定义状态、步进、完成函数 |
| 常用扩展 | httpfs, json, parquet, icu, spellfix | 增强功能 | LOAD 'icu'; -- 启用国际化排序、大小写转换等 | 扩展丰富,可根据需要加载 |
| UDF 限制 | 性能低于内置函数 | 复杂逻辑的补充 | — | 尽量使用内置函数,UDF 用于无法表达的逻辑 |
第9章:Python API 深度使用
聚焦 Python 中的 DuckDB 使用模式。
9.1 duckdb.connect() 参数详解
| 参数 | 语法/值 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| database | 'dbname.db' 或 ':memory:' | 指定数据库文件或内存模式 | # 内存数据库(默认)con = duckdb.connect()# 持久化数据库con = duckdb.connect('analytics.db') | ':memory:' 显式指定内存,None 或 '' 也表示内存 |
| read_only | True / False | 是否以只读模式打开 | con = duckdb.connect('data.db', read_only=True) | 防止意外修改,允许多进程同时读取同一文件 |
| config | {'key': 'value', ...} | 设置运行时配置参数 | con = duckdb.connect(config={'allow_unsigned_extensions': 'true', 'memory_limit': '1GB'}) | 可设置内存、并行度、扩展策略等 |
| check_same_thread | True / False | 是否检查线程一致性 | con = duckdb.connect(check_same_thread=False) | 多线程应用中需设为 False |
| 无参数调用 | duckdb.connect() | 创建临时内存连接 | con = duckdb.connect() # 最常用 | 断开后数据丢失,适合临时分析 |
9.2 cursor 对象方法(execute, fetch_df, fetch_arrow_table 等)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| execute() | con.execute(sql) | 执行 SQL 语句 | con.execute("CREATE TABLE t AS SELECT 42 AS a")result = con.execute("SELECT * FROM t") | 返回结果集对象,可链式调用 fetch 方法 |
| fetchdf() / df() | result.fetchdf() | 获取 Pandas DataFrame | df = con.execute("SELECT * FROM t").df() | fetchdf 和 df 是同义词,最常用 |
| fetcharrowtable() | result.fetcharrowtable() | 获取 PyArrow Table | at = con.execute("SELECT * FROM t").fetcharrowtable() | 高效,零拷贝,适合与 Polars、Dataset 等集成 |
| fetchone() | result.fetchone() | 获取单行元组 | row = con.execute("SELECT a FROM t").fetchone() # (42,) | 用于标量查询或逐行处理 |
| fetchall() | result.fetchall() | 获取所有行的元组列表 | rows = con.execute("SELECT * FROM t").fetchall() # [(42,)] | 全部加载到内存,大数据集慎用 |
| fetchdf_chunked() | result.fetch_record_batch() | 流式获取数据块 | result = con.execute("SELECT * FROM huge_table")while True: batch = result.fetch_record_batch(); if not batch: break; process(batch) | 处理超大数据集,避免内存溢出 |
| close() | result.close() | 关闭结果集 | — | 通常自动管理,显式关闭可释放资源 |
9.3 使用 prepare() 执行预编译语句
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| prepare() | con.prepare(sql) | 预编译带参数的 SQL | prep = con.prepare("SELECT ? + ? AS sum")result = prep.execute([10, 20]).df() | 参数用 ? 占位 |
| 执行预编译语句 | prep.execute([args]) | 使用不同参数多次执行 | prep = con.prepare("INSERT INTO log VALUES (?, ?)")prep.execute(['error', 'File not found'])prep.execute(['info', 'Process started']) | 高效,避免重复解析 SQL |
| 命名参数 | con.prepare("SELECT $1 * $2") | 使用 $1, $2… | prep = con.prepare("SELECT $1 * $2 AS product")prep.execute([5, 6]) | 与位置参数 ? 等效 |
| 性能优势 | 减少解析开销 | 循环中执行相同 SQL 时 | — | 特别适合 ETL 中的批量插入/更新 |
| 生命周期 | prep 对象存活于连接期间 | — | — | 连接关闭后失效 |
9.4 上下文管理器(with 语句)
| 用法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 连接管理 | with duckdb.connect() as con: | 自动关闭连接 | with duckdb.connect('temp.db') as con:con.execute("CREATE TABLE ...")con.execute("INSERT INTO ...")# 连接自动关闭 | 确保资源释放,即使发生异常 |
| 结果集管理 | with con.execute(sql) as res: | 自动关闭结果集 | with con.execute("SELECT * FROM t") as res:df = res.df()# 结果集自动关闭 | 较少使用,通常 fetch 后即释放 |
| 嵌套使用 | with 内创建表/数据 | 临时分析工作流 | with duckdb.connect() as con:con.execute("CREATE VIEW v AS SELECT 1 AS a")df = con.execute("SELECT a+1 FROM v").df() | 内存连接 + with 是安全的临时分析模式 |
| 异常安全 | 自动处理异常 | 防止资源泄漏 | — | 推荐在生产脚本中使用 |
9.5 注册/卸载数据对象(register, unregister)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| register() | con.register('name', obj) | 将 Python 对象注册为 SQL 表 | import pandas as pddf = pd.DataFrame({'x': [1,2]})con.register('my_df', df)result = con.execute("SELECT x*2 FROM my_df").df() | 支持 Pandas DataFrame, PyArrow Table, Polars DataFrame 等 |
| unregister() | con.unregister('name') | 移除注册的表名 | con.unregister('my_df')# my_df 不再可查 | 释放引用,清理命名空间 |
| 作用域 | 仅在当前连接有效 | — | — | 其他连接无法访问 |
| 零拷贝 | 不复制数据 | 高效访问 | — | 特别适合大对象 |
| 动态数据源 | 可注册函数返回表 | 高级用法 | def get_data(): return pd.DataFrame({'val': range(100)})con.register('dynamic', get_data) | 查询时调用函数 |
9.6 与 Pandas 高效交互模式
| 模式 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 注册 DataFrame | 将 df 作为虚拟表查询 | con.register('sales', sales_df)summary = con.execute("SELECT region, SUM(amount) AS total FROM sales GROUP BY region").df() | 避免 pd.read_sql 的低效,利用 DuckDB 引擎处理大 df |
| 执行 SQL 返回 df | 利用 SQL 能力处理 Pandas 数据 | con.register('df', df)cleaned = con.execute("SELECT * FROM df WHERE age BETWEEN 18 AND 65 AND salary IS NOT NULL").df() | 比 Pandas 原生操作更简洁高效 |
| 混合处理 | DuckDB 处理重计算,Pandas 做轻量后处理 | # DuckDB 做聚合agg = con.execute("SELECT cat, AVG(val) FROM df GROUP BY cat").df()# Pandas 做可视化agg.plot(kind='bar') | 发挥各自优势 |
| 大数据集处理 | 用 DuckDB 替代 Pandas 处理大于内存的数据 | # 不要用 pd.read_csv 大文件con.execute("CREATE TABLE t AS SELECT * FROM READ_CSV('huge.csv')")result = con.execute("SELECT ... FROM t ...").df() | 避免 Pandas 内存溢出 |
| 性能对比 | DuckDB 通常比 Pandas 快 | — | — |
第10章:实战应用案例
综合运用所学知识解决实际问题。
10.1 日志分析系统(Nginx 日志解析)
| 步骤 | 方法/语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 日志格式识别 | Nginx 默认 combined 格式 | 了解字段结构 | "$remote_addr - $remote_user [$time_local] \"$request\" $status $body_bytes_sent \"$http_referer\" \"$http_user_agent\"" | 需确认实际使用的日志格式 |
| 读取日志文件 | READ_CSV() + 正则解析 | 加载日志并结构化 | -- 假设日志已按行存为文本CREATE VIEW nginx_logs ASSELECTREGEXP_EXTRACT(line, '(\S+)', 1) AS ip,REGEXP_EXTRACT(line, '\[(.+?)\]', 1) AS timestamp,REGEXP_EXTRACT(line, '"(\S+) (\S+)', 1) AS method,REGEXP_EXTRACT(line, '"(\S+) (\S+)', 2) AS url,CAST(REGEXP_EXTRACT(line, ' (\d{3}) ', 1) AS INTEGER) AS status,CAST(REGEXP_EXTRACT(line, ' (\d+) "$', 1) AS BIGINT) AS bytesFROM READ_CSV('access.log', COLUMNS={'line': 'VARCHAR'}, SEP='\n'); | 使用 SEP='\n' 按行读取,再用正则提取字段 |
| 时间解析 | STRPTIME() | 将字符串转为时间戳 | SELECT STRPTIME(timestamp, '%d/%b/%Y:%H:%M:%S %z') AS ts FROM nginx_logs; | 格式需与日志中 [%d/%b/%Y:%H:%M:%S %z] 一致 |
| 基础统计 | COUNT, GROUP BY | 分析请求量、状态码等 | -- 每小时请求量SELECT DATE_TRUNC('hour', ts) AS hour, COUNT(*) AS reqs FROM nginx_logs GROUP BY hour ORDER BY hour;-- 状态码分布SELECT status, COUNT(*) AS cnt FROM nginx_logs GROUP BY status; | 快速生成监控指标 |
| 用户行为分析 | COUNT(DISTINCT), LIKE | 分析用户、来源、机器人 | -- 独立 IP 数SELECT COUNT(DISTINCT ip) AS unique_ips FROM nginx_logs;-- 来源分析SELECT http_referer, COUNT(*) FROM nginx_logs GROUP BY http_referer;-- 爬虫识别SELECT COUNT(*) FROM nginx_logs WHERE LOWER(user_agent) LIKE '%bot%'; | 结合字符串函数进行模式识别 |
| 性能分析 | AVG, PERCENTILE | 分析响应大小、性能 | SELECT AVG(bytes) AS avg_bytes, PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY bytes) AS p95_bytes FROM nginx_logs; | 识别大文件传输或异常流量 |
| 导出分析结果 | COPY TO | 保存结果供可视化 | COPY (SELECT DATE_TRUNC('day', ts), COUNT(*) FROM nginx_logs GROUP BY 1) TO 'daily_traffic.csv' (HEADER, DELIMITER=','); | 生成报表数据 |
10.2 数据清洗流水线(ETL 简易实现)
| 步骤 | 方法/语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 加载原始数据 | READ_CSV / READ_PARQUET | 读取待清洗数据 | CREATE OR REPLACE VIEW raw_data AS SELECT * FROM READ_CSV('raw_input.csv', AUTO_DETECT=TRUE); | 使用视图隔离原始数据 |
| 处理缺失值 | COALESCE, IS NULL | 填充或过滤空值 | SELECT COALESCE(name, 'Unknown') AS name, age, COALESCE(salary, (SELECT AVG(salary) FROM raw_data)) AS salary FROM raw_data WHERE email IS NOT NULL; -- 过滤关键字段为空 | 根据业务逻辑决定填充策略 |
| 数据类型修正 | CAST, TRY_CAST | 确保类型正确 | SELECT name, TRY_CAST(age_str AS INTEGER) AS age, TRY_CAST(price_txt AS DOUBLE) AS price FROM raw_data WHERE TRY_CAST(age_str AS INTEGER) IS NOT NULL; -- 过滤转换失败 | TRY_CAST 避免因脏数据报错 |
| 去重 | ROW_NUMBER() 或 DISTINCT | 移除重复记录 | -- 方法1: 使用窗口函数保留第一条CREATE VIEW cleaned_1 ASSELECT * EXCEPT(rn) FROM (SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY updated_at DESC) AS rn FROM raw_data) WHERE rn = 1;-- 方法2: 简单去重SELECT DISTINCT * FROM raw_data; | 根据主键和时间决定保留策略 |
| 标准化文本 | TRIM, UPPER, REPLACE | 统一文本格式 | SELECT TRIM(UPPER(name)) AS name, REPLACE(phone, '-', '') AS phone_clean FROM cleaned_1; | 确保一致性,便于后续匹配 |
| 创建维度表 | CREATE TABLE AS SELECT | 构建规范化表 | CREATE TABLE dim_product AS SELECT DISTINCT product_id, product_name, category FROM cleaned_data;CREATE TABLE fact_sales AS SELECT sale_id, product_id, amount, sale_date FROM cleaned_data; | 实现简单的维度建模 |
| 导出清洗后数据 | COPY TO | 保存结果 | COPY fact_sales TO 'cleaned_sales.parquet' (FORMAT PARQUET);COPY dim_product TO 'dim_product.parquet' (FORMAT PARQUET); | 推荐使用 Parquet 作为中间/最终存储 |
10.3 本地 OLAP 分析(销售数据分析)
| 分析类型 | SQL 查询 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 基础聚合 | SUM, COUNT, GROUP BY | 总体销售情况 | SELECT SUM(amount) AS total_sales, COUNT(*) AS order_count, AVG(amount) AS avg_order_value FROM sales; | 快速了解业务规模 |
| 时间趋势分析 | DATE_TRUNC, GROUP BY | 销售随时间变化 | SELECT DATE_TRUNC('month', order_date) AS month, SUM(amount) AS monthly_sales FROM sales GROUP BY month ORDER BY month; | 识别季节性、增长趋势 |
| 产品分析 | GROUP BY, ORDER BY | 畅销产品排行 | SELECT product_name, SUM(amount) AS product_sales, COUNT(*) AS units_sold FROM sales GROUP BY product_name ORDER BY product_sales DESC LIMIT 10; | 识别核心产品 |
| 区域分析 | GROUP BY, JOIN | 各地区业绩对比 | SELECT r.region_name, SUM(s.amount) AS region_sales FROM sales s JOIN regions r ON s.region_id = r.id GROUP BY r.region_name; | 结合地理维度 |
| 客户分层 | CASE, SUM | RFM 或自定义分层 | SELECT customer_id, SUM(amount) AS total_spent, CASE WHEN SUM(amount) > 10000 THEN 'VIP' WHEN SUM(amount) > 1000 THEN 'Regular' ELSE 'New' END AS customer_tier FROM sales GROUP BY customer_id; | 支持精细化运营 |
| 同比/环比 | LAG(), DATE arithmetic | 增长率计算 | WITH monthly AS (SELECT DATE_TRUNC('month', order_date) AS m, SUM(amount) AS sales FROM sales GROUP BY m) SELECT m, sales, LAG(sales, 1) OVER (ORDER BY m) AS prev_month, (sales - LAG(sales, 1) OVER (ORDER BY m)) / LAG(sales, 1) OVER (ORDER BY m) AS growth_rate FROM monthly; | 分析业务健康度 |
| 多维分析(CUBE) | GROUPING SETS / CUBE | 全组合聚合 | SELECT region, product_category, SUM(amount) AS sales FROM sales GROUP BY CUBE(region, product_category); | 生成交叉报表,但可能数据量大 |
10.4 嵌入式分析服务(FastAPI + DuckDB)
| 组件/步骤 | 实现方式 | 用途 | 代码示例(Python) | 注意事项 |
|---|---|---|---|---|
| 项目结构 | FastAPI + duckdb | 创建 REST API | from fastapi import FastAPIimport duckdbapp = FastAPI()con = duckdb.connect('analytics.db', read_only=True) # 只读连接 | 将 DuckDB 作为嵌入式分析引擎 |
| 健康检查 | @app.get("/health") | 服务可用性检测 | @app.get("/health")def health_check():return {"status": "ok"} | 标准运维接口 |
| 参数化查询 API | @app.get("/sales") | 提供数据查询接口 | @app.get("/sales")def get_sales(region: str = None, start_date: str = None):query = "SELECT region, SUM(amount) FROM sales WHERE 1=1"params = []if region: query += " AND region = ?"; params.append(region)if start_date: query += " AND order_date >= ?"; params.append(start_date)query += " GROUP BY region"df = con.execute(query, params).df()return df.to_dict(orient='records') | 使用预编译语句防止 SQL 注入 |
| 返回 JSON 结果 | .df().to_dict() | 将结果转为 JSON 响应 | # 在路由函数中return df.to_dict(orient='records') | FastAPI 自动序列化 |
| 静态文件服务 | @app.get("/") | 提供前端页面 | from fastapi.staticfiles import StaticFilesapp.mount("/", StaticFiles(directory="frontend", html=True), name="frontend") | 构建完整 Web 应用 |
| 连接池管理 | 全局连接或依赖注入 | 管理数据库连接 | # 简单场景:全局只读连接# 复杂场景:使用 Depends 获取连接 | 高并发时考虑连接池或每个请求新连接 |
| 部署 | uvicorn.run() | 启动服务 | if __name__ == "__main__":import uvicornuvicorn.run(app, host="0.0.0.0", port=8000) | 可打包为 Docker 镜像或独立服务 |
总结: DuckDB 凭借其轻量、高性能、零依赖的特性,非常适合嵌入到 Python 应用中,为 Web 服务、数据分析工具、桌面应用等提供强大的本地分析能力。通过与 FastAPI 等框架结合,可以快速构建出功能完整的嵌入式 BI 或数据服务。