本系列基于 SQLMesh 官方文档(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)整理,共 3 篇,面向初学者。本篇是第一篇,带你搞清楚"Python 模型是什么、怎么定义、必填项有哪些"。
- 第(二)篇:取数与依赖管理、四种引擎的 DataFrame 实战
- 第(三)篇:前后置语句、蓝图批量建模、避坑清单
1. 为什么需要 Python 模型?
在数据管道里,SQL 是最常用的工具,但有些场景 SQL 表达起来很吃力:
- 机器学习管道:需要调用训练好的模型做预测;
- 与外部 API 交互:比如拉取第三方服务的实时数据;
- 复杂的业务逻辑:多层条件分支、循环、解析嵌套结构等,用 SQL 写会非常痛苦。
SQLMesh 对 Python 模型提供了一等公民(first-class)支持:只要你的函数最终返回一个 Pandas 或 Spark(以及 Snowpark、Bigframe)DataFrame,模型里几乎可以做任何事情。
不过要注意,Python 模型不支持以下几种模型 kind,需要用这些 kind 时请改用 SQL 模型:
| 不支持的 kind | 说明 |
|---|---|
VIEW | 视图模型 |
SEED | 种子(静态 CSV 加载)模型 |
MANAGED | 托管表模型 |
EMBEDDED | 嵌入式模型 |
提示:SQL 模型的默认 kind 是
VIEW,而Python 模型的默认 kind 是FULL——即每次评估时重新计算并写一张全量表。
2. 模型的定义:文件放哪、函数怎么写
创建 Python 模型的步骤只有两步:
- 在 SQLMesh 项目的
models/目录下新建一个*.py文件; - 在文件里定义一个名为
execute的函数,并用@model装饰器包裹。
下面是最小可用骨架:
importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,model@model("my_model.name",# 模型名,通常是 "schema.表名" 的形式columns={# 【必填】输出 DataFrame 的列名 -> 类型"column_name":"int",},)defexecute(context:ExecutionContext,# 执行上下文:跑查询、拿环境变量start:datetime,# 本次处理时间区间的起点end:datetime,# 本次处理时间区间的终点execution_time:datetime,# 实际执行时刻**kwargs:t.Any,# 运行时传入的任意键值参数)->pd.DataFrame:returnpd.DataFrame({"column_name":[1]})逐个拆解关键点
①@model装饰器 = SQL 模型里的MODELDDL
它负责声明模型的元数据(名字、kind、列、调度等)。参数命名与 SQL 模型MODELDDL 中的字段完全一致,所以会写 SQL 模型的人可以无缝迁移。
②columns为什么是必填的?
SQLMesh 会在评估模型之前先在引擎里建好表,因此必须提前知道输出数据的 schema(列名和类型)。这与 SQL 模型不同——SQL 模型的列名和类型可以由 SQLMesh 从查询里自动推断,Python 代码无法静态解析,所以必须手动声明。
⚠️ 大坑预警:如果实际返回的 DataFrame 与
columns声明不一致(多了列、少了列、类型不符),会产生意外行为甚至报错。声明时就对齐,是初学者的第一守则。
③execute函数的参数
context(ExecutionContext)是最核心的参数:能执行 SQL 查询、获取当前处理的时间区间、读取变量;start/end:对增量模型来说,代表本次要处理的时间窗口;**kwargs:接收运行时传入的额外参数。
④ 返回值
可以返回Pandas、PySpark、Bigframe 或 Snowpark的 DataFrame 实例。如果输出数据量太大,还可以用Python 生成器(generator)分块返回(详见第(二)篇)。
3. 第一个完整可运行示例
展示元数据常用字段的完整用法(字段名与 SQLMODELDDL 一致):
importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlglot.expressionsimportto_columnfromsqlmeshimportExecutionContext,model@model("docs_example.basic",owner="janet",cron="@daily",columns={"id":"int","name":"text",},column_descriptions={"id":"Unique ID","name":"Name corresponding to the ID",},audits=[("not_null",{"columns":[to_column("id")]}),],)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:returnpd.DataFrame([{"id":1,"name":"name"}])这个例子返回一个静态 DataFrame,没有实际业务意义,但它把owner(负责人)、cron(调度)、column_descriptions(列注释)、audits(质量校验)等常用元数据都示范了一遍。
补充:Python 模型的列注释不能像 SQL 模型那样从代码行内注释推断,必须写在
@model的column_descriptions里;若其中出现了columns里没有的列名,SQLMesh 会直接报错。
4. 配置模型 kind:以增量模型为例
模型 kind 决定它如何被计算和存储。Python 模型里用字典来声明 kind,name键的值必须是ModelKindName枚举的成员,且该枚举需要在文件开头先导入。
可支持的name取值包括:
ModelKindName.FULL(默认,全量重算)ModelKindName.INCREMENTAL_BY_TIME_RANGEModelKindName.INCREMENTAL_BY_UNIQUE_KEYModelKindName.INCREMENTAL_BY_PARTITIONModelKindName.SCD_TYPE_2_BY_TIME/SCD_TYPE_2_BY_COLUMNModelKindName.CUSTOM/EXTERNAL
最常用的"增量按时间区间"写法:
fromsqlmeshimportExecutionContext,modelfromsqlmesh.core.model.kindimportModelKindName@model("docs_example.incremental_model",kind=dict(name=ModelKindName.INCREMENTAL_BY_TIME_RANGE,time_column="model_time_column",# 用哪一列做时间分区/覆盖判断),)defexecute(context,start,end,execution_time,**kwargs):...配置 kind 后,SQLMesh 每次只会处理start~end之间的数据,而不是全量重算——这是生产环境管道的标配。
5. 执行上下文:取数的第一步
Python 模型可以随心所欲地做任何事,但官方强烈建议所有模型保持幂等(idempotent):同一个区间重跑多次,结果应该一致。
从上游取数最简单的方式是fetchdf,它接收一段 SQL 并返回 DataFrame:
df=context.fetchdf("SELECT * FROM my_table")注意:fetchdf返回的类型跟随执行引擎——引擎是 Spark 时返回的就是 Spark DataFrame。
6. 读取用户自定义变量
项目配置里定义的全局变量,在 Python 模型中有两种访问方式。
方式一:context.var(name, default)
@model("my_model.name")defexecute(context,start,end,execution_time,**kwargs):var_value=context.var("var")var_with_default_value=context.var("var_with_default","default_value")...方式二:直接作为execute的函数参数(参数名 = 变量名)
@model("my_model.name")defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,my_var:Optional[str]=None,# 必须给默认值,防止变量缺失**kwargs:t.Any,):my_var_plus1=my_var+1...两个注意事项:
- 参数必须显式声明——变量不能通过
kwargs偷偷获取; - 变量可能不存在时,必须给默认值,否则运行时会出问题。
小结
本篇覆盖了 Python 模型的地基知识:
- 一个
@model装饰器 + 一个execute函数 = 一个模型,元数据字段与 SQL 模型一一对应; - schema 先于代码——
columns必填且必须与返回的 DataFrame 严格一致; - kind 用字典 +
ModelKindName枚举声明,增量生产用INCREMENTAL_BY_TIME_RANGE; - 取数靠
context.fetchdf,变量靠context.var或显式函数参数。
下一篇我们将走进真实的数据流转:如何声明上游依赖、如何用 PySpark/Snowpark/Bigframe 让计算下推到集群,以及大输出分块与空表处理的正确姿势。
参考资料:SQLMesh 官方文档 — Python models(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)