☰
SQLMesh Python 模型入门(一):基础语法与核心概念
2026/10/3 7:16:22 网站建设 项目流程

本系列基于 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 模型的步骤只有两步:

  1. 在 SQLMesh 项目的models/目录下新建一个*.py文件;
  2. 在文件里定义一个名为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_RANGE
  • ModelKindName.INCREMENTAL_BY_UNIQUE_KEY
  • ModelKindName.INCREMENTAL_BY_PARTITION
  • ModelKindName.SCD_TYPE_2_BY_TIME/SCD_TYPE_2_BY_COLUMN
  • ModelKindName.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 模型的地基知识:

  1. 一个@model装饰器 + 一个execute函数 = 一个模型,元数据字段与 SQL 模型一一对应;
  2. schema 先于代码——columns必填且必须与返回的 DataFrame 严格一致;
  3. kind 用字典 +ModelKindName枚举声明,增量生产用INCREMENTAL_BY_TIME_RANGE;
  4. 取数靠context.fetchdf,变量靠context.var或显式函数参数。

下一篇我们将走进真实的数据流转:如何声明上游依赖、如何用 PySpark/Snowpark/Bigframe 让计算下推到集群,以及大输出分块与空表处理的正确姿势。


参考资料:SQLMesh 官方文档 — Python models(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询