Daft SQL SELECT 语句实战指南:从基础查询到引擎级执行原理
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
Daft 为 AI 与多模态数据场景设计的高性能数据引擎,其内置 SQL 方言紧密贴近 DuckDB 与 PostgreSQL(见 docs/sql/index.md),而SELECT语句正是所有查询的入口。本文以仓库 SQL 参考文档 docs/sql/statements/select.md 为骨架,完整覆盖其全部示例,并深入daft.sqlPython API、daft-sql规划器源码与tests/sql测试套件,帮助你从"会写"到"懂原理",掌握 Daft 中 SELECT 的完整用法、执行流程与能力边界。
SELECT 语句是什么
在 Daft 中,SELECT语句用于查询某个 catalog(数据目录)中的表,它同时支持:
- 直接求值表达式:无需任何表,如
SELECT 1 + 1; - 查询 DataFrame:
daft.sql()会自动把当前 Python 作用域中的daft.DataFrame变量注册为可查询的表 - 查询外部数据源:通过
read_parquet、read_csv、read_iceberg等表函数直接读取文件与湖格式
从源码看,Daft 的 SQL 解析与规划基于sqlparsercrate,在 src/daft-sql/src/planner.rs 中实现;顶层语句类型定义在 src/daft-sql/src/statement.rs,其中Statement::Select(Select)直接承载SELECT ...查询,Select本质就是一个LogicalPlanRef——即 SELECT 查询会被翻译成 Daft 的逻辑计划,最终与 DataFrame API 走同一条执行链路。
运行环境准备:三种执行入口
在写 SELECT 之前,先明确 Daft 提供哪几种执行入口:
1. 模块级函数daft.sql(sql, ...)
定义于 daft/sql/sql.py,是使用最频繁的入口:
import daft df1 = daft.from_pydict({"a": [1, 2, 3], "b": ["foo", "bar", "baz"]}) df2 = daft.from_pydict({"a": [1, 2, 3], "c": ["daft", None, None]}) # Daft 自动从 Python 全局命名空间识别 df1 和 df2 result_df = daft.sql("SELECT * FROM df1 JOIN df2 ON df1.a = df2.a") result_df.show()关键参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
sql | str | 必填 | 要执行的 SQL 查询 |
register_globals | bool | True | 是否把调用方作用域中的 DataFrame 变量注册进 catalog。关闭后需通过bindings显式传入,否则报错(见 tests/sql/test_sql.py 中test_sql_function_register_globals) |
**bindings | DataFrame | 无 | 额外的 DataFrame 绑定(可视为 CTE 表),优先级最高,可覆盖同名全局变量 |
从实现看,daft.sql内部按"全局变量 → bindings(CTE)"的顺序构建表绑定字典(py_ctes),显式传入的bindings放在最后所以不会被遮蔽(对应 daft/sql/sql.py);随后调用底层_sql_exec执行,并把结果包装成新的DataFrame,无缝融入 Daft 的懒执行(lazy)计划。
2. Session 级方法sess.sql(...)
如果你使用Session管理 catalog,则通过会话执行,并可先用create_temp_table注册临时表(示例见 docs/sql/index.md):
from daft import Session sess = Session() sess.create_temp_table("T", daft.from_pydict({"a": [0, 1]})) sess.create_temp_table("S", daft.from_pydict({"b": [1, 0]})) sess.sql("SELECT * FROM T, S").show()关于 Session 与 catalog 连接的详细用法,参见 docs/configuration/sessions-usage.md。
3. 表达式级函数daft.sql_expr(sql)
daft.sql_expr把一段 SQL 表达式直接解析为Expression(daft/sql/sql.py),可嵌入 DataFrame 操作:
df = daft.from_pydict({"a": [1, 2, 3], "b": [4, 5, 6]}) df = df.with_column("c", daft.sql_expr("a + b")) # 与 col("a") + col("b") 等价 df.show()它也会在部分 DataFrame 操作中被自动调用,例如df.where("x < 3 AND y > 4")里的字符串过滤条件(daft/sql/sql.py)。
基础 SELECT 用法(原文档示例全集)
下面完整列出 docs/sql/statements/select.md 中的全部示例,并补充可运行的上下文。
求值单个表达式
不依赖任何表,直接求值常量表达式:
SELECT 1 + 1;在 Daft 中执行等价于:
daft.sql("SELECT 1 + 1").show()选择全部列
从表T中取出所有列与所有行:
SELECT * FROM T;选择指定列
从表T中选出a、b、c三列:
SELECT a, b, c FROM T;投影列的顺序即输出 DataFrame 的列顺序。
在投影中应用标量函数
对列a、b分别应用标量函数foo和bar:
SELECT foo(a), bar(b) FROM T;Daft 的 SQL 函数集非常庞大——聚合、字符串、时间、列表、URI、图像处理等函数均可在 SELECT 投影中使用,详见下文"函数调用"一节。
计数非空值
统计列a非空的行数:
SELECT COUNT(a) FROM T;注意:COUNT(a)只统计a非NULL的行;若要统计所有行,用COUNT(*)。测试 tests/sql/test_sql.py 的test_sql_count_star同时验证了两种写法。
分组计数
按列b分组,统计每组行数:
SELECT COUNT(*), b FROM T GROUP BY b;这是最基础的GROUP BY聚合查询,其执行路径会进入规划器中的聚合分支(见下节源码解析)。
引擎级原理:SELECT 的处理顺序
搞清楚 SELECT 的底层处理顺序,能帮你避开大量 SQL 陷阱。从 src/daft-sql/src/planner.rs 的plan_query可以看出,Daft 严格按照以下顺序规划一条 SELECT 查询:
- CTE 绑定:解析
WITH子句,将公共表表达式注册进PlannerContext(plan_ctes) - FROM / JOIN:解析数据来源与连接(
plan_from) - SELECT 投影:把投影列表中的每一项翻译为表达式(
select_item_to_expr) - WHERE:解析过滤谓词并施加
plan.filter(filter) - GROUP BY:解析分组表达式(支持表达式分组、
ROLLUP、按投影序号分组) - ORDER BY:解析排序键(
plan_order_by_exprs) - 聚合判定:若投影中包含聚合函数或存在 GROUP BY,走
plan_aggregate_query(此时解析HAVING);否则走plan_non_agg_query - DISTINCT:
SELECT DISTINCT或DISTINCT ON (cols)施加plan.distinct(...) - LIMIT / OFFSET:最后施加分页
值得注意的两个细节:
- OFFSET 与 LIMIT 的父子关系:规划器中有一条注释明确说明:由于 Daft SQL 方言紧跟 DuckDB 与 PostgreSQL,当
LIMIT与OFFSET同时出现时,无论书写顺序如何,都会保证 OFFSET 是 LIMIT 的子节点(src/daft-sql/src/planner.rs)。LIMIT 7 OFFSET 2与OFFSET 2 LIMIT 7语义一致。 - 取值校验:
LIMIT n与OFFSET n都必须是非负整数常量,否则直接报错——测试 tests/sql/test_limit_offset.py 中test_negative_limit断言错误信息"LIMIT <n> must be greater than or equal to 0, instead got: -1",test_negative_offset同理。另外OFFSET 必须与 LIMIT 搭配使用,单独OFFSET 17会抛出"Offset without limit is unsupported now!"(见test_offset_without_limit)。
此外,规划器在语句级还会对能力边界做前置检查:Subqueries are not supported(FROM 中的子查询不被支持)、VALUES are not supported、INSERT/UPDATE/DELETE/MERGE均不支持(src/daft-sql/src/planner.rs),遇到会抛出带^定位符的友好错误。
聚合与 GROUP BY 的完整用法
原文档给出了COUNT两种形态,而 Daft 的 SQL 聚合远不止于此。tests/sql/test_aggs.py 的test_aggs_sql一次验证了十余种聚合函数在 SQL 与 DataFrame API 下结果完全一致:
SELECT sum(values) as sum, product(values) as product, mean(values) as mean, avg(values) as avg, percentile(values, 0.99) as p99, median(values) as median, min(values) as min, max(values) as max, count(values) as count, count(distinct values) as count_distinct, stddev(values) as std, stddev_pop(values) AS std_pop, stddev_samp(values) AS std_samp, variance(values) AS variance, var(values) AS var, var_samp(values) AS var_samp, var_pop(values) AS var_pop FROM df常用聚合函数速查:
| 函数 | 语义 |
|---|---|
COUNT(expr)/COUNT(*) | 非空计数 / 全行计数 |
COUNT(DISTINCT expr) | 去重计数 |
SUM/PRODUCT | 求和 / 求积 |
AVG/MEAN | 平均值 |
MIN/MAX | 最小值 / 最大值 |
MEDIAN/PERCENTILE(expr, p) | 中位数 / 分位数(如percentile(values, 0.99)表示 P99) |
STDDEV/STDDEV_POP/STDDEV_SAMP | 标准差(总体 / 样本) |
VARIANCE/VAR_POP/VAR_SAMP | 方差(总体 / 样本) |
HAVING 过滤分组
与GROUP BY配合,HAVING对聚合后的分组结果过滤,测试覆盖在 tests/sql/test_aggs.py 的test_having系列:
SELECT b, SUM(c) AS total FROM T GROUP BY b HAVING SUM(c) > 100;按投影序号分组与 ROLLUP
规划器支持按 SELECT 投影中的序号分组(即 DuckDB/PostgreSQL 风格的GROUP BY 1, 2),对应测试test_group_by_ordinal_1/2/1_2、test_group_by_ordinal_zero_raises(序号从 0 开始会报错)等;同时支持ROLLUP生成小计行,如 tests/sql/test_aggs.py 中test_simple_rollup对GROUP BY ROLLUP(dept, year)的输出会额外产生dept/year为None的小计行。
函数调用:SQL 嵌套等价于 Python 方法链
Daft 的 SQL 可调用全部Expression能力。与 Python API 的方法链风格(col("a").download().decode_image())不同,SQL 中需要改用函数嵌套:image_decode(url_download(a))(示例见 docs/sql/index.md 的 SQL Functions 一节)。
df = daft.from_pydict({"urls": [ "https://user-images.githubusercontent.com/17691182/190476440-28f29e87-8e3b-41c4-9c28-e112e595f558.png", # ...更多 URL ]}) # SQL 版本:函数嵌套 df = daft.sql("SELECT image_decode(url_download(urls)) FROM df") # 等价的 Python 版本:方法链 df = df.select(daft.col("urls").download().decode_image())两者输出相同的Image[MIXED]类型列。这也意味着你可以在 SQL 中直接完成多模态数据流水线(下载 → 解码 → 后续图像处理),而不是只做关系型查询。
更妙的是表达式级互通:daft.sql_expr("A + B as C")与(daft.col("A") + daft.col("B")).alias("C")打印出来完全一致,均为col(A) + col(B) as C(docs/sql/index.md),说明 SQL 表达式与 Python 表达式在内部是同一套 DSL,可以在一个 Pipeline 里自由混用两种语法。
排序与分页:ORDER BY / LIMIT / OFFSET
虽然原文档未展开,但排序与分页是 SELECT 的高频配套子句,且 Daft 行为有明确规范(测试全量覆盖见 tests/sql/test_limit_offset.py):
-- 按 id 升序取前 7 行 SELECT name FROM input_df ORDER BY id LIMIT 7; -- 降序 + 跳过 2 行 + 取 7 行(OFFSET 写在 LIMIT 前后均可) SELECT id, name FROM input_df ORDER BY id DESC OFFSET 2 LIMIT 7; -- 经典分页 SELECT id, name FROM input_df ORDER BY id OFFSET {offset} LIMIT {limit};要点归纳:
LIMIT/OFFSET必须是>= 0的整数常量,负值直接报错;OFFSET不能脱离LIMIT单独使用;LIMIT与OFFSET顺序无关,规划器统一保证 OFFSET 在 LIMIT 之下;- 可与
ORDER BY、WHERE、JOIN自由组合,且支持在子查询层叠使用(见test_paging、test_offset_limit_with_join); - 顶层
ORDER BY需要配合确定性的排序键,才能保证分页结果稳定。
另外 Daft 还支持SELECT DISTINCT与DISTINCT ON (columns)(见规划器 src/daft-sql/src/planner.rs)。
集合操作:UNION / INTERSECT
在plan_query的集合操作分支中(src/daft-sql/src/planner.rs),支持:
SELECT a FROM T1 UNION ALL SELECT a FROM T2; SELECT a FROM T1 UNION SELECT a FROM T2; -- 隐式去重(UNION DISTINCT) SELECT a FROM T1 INTERSECT ALL SELECT a FROM T2; -- 保留重复 INTERSECT SELECT a FROM T2; -- 去重其中UNION/UNION DISTINCT默认按位置合并,也支持BY NAME变体(UNION ALL BY NAME等)按列名合并;对应测试见 tests/sql/set_ops.py。EXCEPT目前不在支持列表内。
用表函数直接读取数据源
SELECT 的FROM不仅能接表,还能接表函数直接读文件。read_parquet、read_csv、read_json、read_iceberg、read_deltalake的选项完整清单见 docs/sql/index.md 的 Table Function Options 表格。几个实用模式:
# 路径 + 命名参数(=> 或 := 均可) daft.sql("SELECT * FROM read_csv('data.csv', has_headers => false)") daft.sql("SELECT * FROM read_csv('data.csv', has_headers := false)") # path 也可作为命名参数 daft.sql("SELECT * FROM read_csv(path => 'data.csv')") # 多文件用 SQL 数组 daft.sql("SELECT * FROM read_parquet(['a.parquet', 'b.parquet'])") # 显式 schema:struct 字面量 daft.sql("""SELECT * FROM read_csv('data.csv', schema := {'a': 'int64', 'b': 'string'})""") # Iceberg 分支读取(snapshot_id / branch / tag 互斥) daft.sql("SELECT * FROM read_iceberg('/warehouse/db/t/metadata/v3.metadata.json', branch => 'audit')") # 跳过损坏文件:collect 后通过 df.skipped_corrupt_files 查看 df = daft.sql("SELECT * FROM read_csv('s3://my-bucket/data/**/*.csv', ignore_corrupt_files => true)") df.collect() print(df.skipped_corrupt_files)以read_parquet为例,其可用选项包括infer_schema、schema、coerce_int96_timestamp_unit、chunk_size、multithreaded、io_config、file_path_column、hive_partitioning、ignore_corrupt_files,与 Python 侧daft.read_parquet一一对应;ignore_corrupt_files的深入说明见 docs/connectors/generic-file-source-options.md。
窗口函数:SELECT 投影中的进阶分析
Daft SQL 支持在 SELECT 投影中使用窗口函数,语法为:
function_name([expr]) OVER ( [PARTITION BY expr_list] [ORDER BY order_list] [frame_clause] )支持的函数分三类(完整文档见 docs/sql/window_functions.md):
- 排名函数:
ROW_NUMBER()、RANK()(并列留空位)、DENSE_RANK()(并列不留空位) - 偏移函数:
LAG(value [, offset [, default]])、LEAD(value [, offset [, default]]),offset 缺省为 1 - 聚合函数:所有聚合函数都可用作窗口函数,如
SUM/AVG/COUNT/MIN/MAX
-- 组内排名 SELECT category, value, ROW_NUMBER() OVER (PARTITION BY category ORDER BY value) AS row_num, RANK() OVER (PARTITION BY category ORDER BY value) AS rank FROM sales; -- 组内累计(默认帧:UNBOUNDED PRECEDING 到 CURRENT ROW) SELECT category, value, SUM(value) OVER (PARTITION BY category ORDER BY value) AS running_sum FROM sales; -- 滑动平均(当前行 + 前 2 行) SELECT date, value, AVG(value) OVER (ORDER BY date ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) AS moving_avg FROM time_series;窗口帧支持ROWS模式(UNBOUNDED PRECEDING/n PRECEDING/CURRENT ROW/n FOLLOWING/UNBOUNDED FOLLOWING),RANGE模式尚未完全支持;同时存在以下限制:全局分区(无PARTITION BY)、WINDOW命名子句、IGNORE/RESPECT NULLS均暂不支持。
错误信息与常见坑
Daft 的 SQL 解析错误会给出带^定位符的友好提示:规划器从sqlparser错误中提取行列号,并把错误位置直接标注在原始 SQL 文本上(src/daft-sql/src/planner.rs),对应测试test_sql_caret_error_eof、test_sql_caret_error_multiline等(tests/sql/test_sql.py)。
实操中最容易踩的坑汇总:
| 情况 | 结果 |
|---|---|
LIMIT -1/OFFSET -1 | 报错:必须>= 0 |
单独OFFSET n(无 LIMIT) | 报错:Offset without limit is unsupported |
| FROM 中出现子查询 | 报错:Subqueries are not supported |
多条语句SELECT ...; SELECT ... | 报错:多语句不支持(见test_sql_multi_statement_sql_error) |
表名与关键字冲突(如TABLE) | 需转义处理,测试见test_sql_function_table_name_is_keyword |
GROUP BY 0或超范围序号 | 报错(见test_group_by_ordinal_zero_raises) |
关于标识符大小写、引号与命名规范,参见 docs/sql/identifiers.md;数据类型体系参见 docs/sql/datatypes.md。
结语与当前状态
综上所述,Daft 的 SELECT 语句虽以"查询 catalog 中的表"为起点,实际上已经成为融合关系查询、聚合分析、窗口计算、多模态函数、外部数据源读取与 DataFrame 生态的统一入口。官方在 docs/sql/statements/select.md 中明确标注 SQL Reference 文档仍在完善中(Work in Progress),源码中daft.sql的 docstring 也提示该功能"早期开发中,API 可能变动"——因此建议以当前仓库版本为准,并结合 tests/sql 目录下的测试样例(test_sql.py、test_aggs.py、test_limit_offset.py、test_joins.py、test_window.py等)验证具体行为。掌握本文的基础示例与规划器执行顺序,你就具备了在 Daft 中编写、调试与优化 SELECT 查询的完整能力。
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考