资讯动态

SQLMesh Python 模型入门(一):基础语法与核心概念

发布时间:2026/10/3 7:16:35 来源:尧图企业网站定制
本系列基于 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、BigframeDataFrame模型里几乎可以做任何事情。不过要注意Python 模型不支持以下几种模型 kind需要用这些 kind 时请改用 SQL 模型不支持的 kind说明VIEW视图模型SEED种子静态 CSV 加载模型MANAGED托管表模型EMBEDDED嵌入式模型提示SQL 模型的默认 kind 是VIEW而Python 模型的默认 kind 是FULL——即每次评估时重新计算并写一张全量表。2. 模型的定义文件放哪、函数怎么写创建 Python 模型的步骤只有两步在 SQLMesh 项目的models/目录下新建一个*.py文件在文件里定义一个名为execute的函数并用model装饰器包裹。下面是最小可用骨架importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,modelmodel(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函数的参数contextExecutionContext是最核心的参数能执行 SQL 查询、获取当前处理的时间区间、读取变量start/end对增量模型来说代表本次要处理的时间窗口**kwargs接收运行时传入的额外参数。④ 返回值可以返回Pandas、PySpark、Bigframe 或 Snowpark的 DataFrame 实例。如果输出数据量太大还可以用Python 生成器generator分块返回详见第二篇。3. 第一个完整可运行示例展示元数据常用字段的完整用法字段名与 SQLMODELDDL 一致importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlglot.expressionsimportto_columnfromsqlmeshimportExecutionContext,modelmodel(docs_example.basic,ownerjanet,crondaily,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 模型里用字典来声明 kindname键的值必须是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.kindimportModelKindNamemodel(docs_example.incremental_model,kinddict(nameModelKindName.INCREMENTAL_BY_TIME_RANGE,time_columnmodel_time_column,# 用哪一列做时间分区/覆盖判断),)defexecute(context,start,end,execution_time,**kwargs):...配置 kind 后SQLMesh 每次只会处理start~end之间的数据而不是全量重算——这是生产环境管道的标配。5. 执行上下文取数的第一步Python 模型可以随心所欲地做任何事但官方强烈建议所有模型保持幂等idempotent同一个区间重跑多次结果应该一致。从上游取数最简单的方式是fetchdf它接收一段 SQL 并返回 DataFramedfcontext.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_valuecontext.var(var)var_with_default_valuecontext.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_plus1my_var1...两个注意事项参数必须显式声明——变量不能通过kwargs偷偷获取变量可能不存在时必须给默认值否则运行时会出问题。小结本篇覆盖了 Python 模型的地基知识一个model装饰器 一个execute函数 一个模型元数据字段与 SQL 模型一一对应schema 先于代码——columns必填且必须与返回的 DataFrame 严格一致kind 用字典 ModelKindName枚举声明增量生产用INCREMENTAL_BY_TIME_RANGE取数靠context.fetchdf变量靠context.var或显式函数参数。下一篇我们将走进真实的数据流转如何声明上游依赖、如何用 PySpark/Snowpark/Bigframe 让计算下推到集群以及大输出分块与空表处理的正确姿势。参考资料SQLMesh 官方文档 — Python modelshttps://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑