本系列基于 SQLMesh 官方文档(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)整理,共 3 篇,面向初学者。本篇是完结篇,覆盖工程化进阶能力,并给出全系列的避坑总表。
- 第(一)篇:基础语法与核心概念
- 第(二)篇:取数与依赖管理、四种引擎的 DataFrame 实战
1. 前后置语句(pre/post-statements)
前置/后置语句让你在模型运行前后执行 SQL。典型用途:修改会话设置、创建索引。
⚠️ 并发提醒:不要写会与其他并发模型冲突的语句(例如创建物理表),并发执行时行为不可预测。
1.1 在装饰器里声明
pre_statements/post_statements接收一个列表,元素可以是 SQL 字符串、SQLGlot 表达式或宏调用:
@model("db.test_model",kind="full",columns={"id":"int","name":"text",},pre_statements=["SET GLOBAL parameter = 'value';",exp.Cache(this=exp.table_("x"),expression=exp.select("1")),],post_statements=["@CREATE_INDEX(@this_model, id)"],)defexecute(context,start,end,execution_time,**kwargs)->pd.DataFrame:returnpd.DataFrame([{"id":1,"name":"name"}])其中@CREATE_INDEX是自定义宏,在项目的macros目录里这样定义——仅在creating(建表)阶段执行:
@macro()defcreate_index(evaluator:MacroEvaluator,model_name:str,column:str,):ifevaluator.runtime_stage=="creating":returnf"CREATE INDEX idx ON{model_name}({column});"returnNone项目级默认值:也可以在配置的model_defaults里为整个项目定义 pre/post 语句,所有模型自动继承,并与模型级语句合并(默认语句先执行)。
1.2 在函数体内声明
规则很简单:
- 前置语句:写在
return/yield之前任意位置即可; - 后置语句:必须把
return改成yield,然后写在yield之后(因为后置语句要在函数产出数据之后才执行)。
defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:# pre-statementcontext.engine_adapter.execute("SET GLOBAL parameter = 'value';")# post-statement 要求用 yield 而不是 returnyieldpd.DataFrame([{"id":1,"name":"name"}])# post-statementcontext.engine_adapter.execute("CREATE INDEX idx ON example.pre_post_statements (id);")2. on_virtual_update:虚拟层更新后执行
on_virtual_update在 Virtual Update 完成后执行 SQL,典型用途是给虚拟层的视图授权:
@model("db.test_model",kind="full",columns={"id":"int","name":"text",},on_virtual_update=["GRANT SELECT ON VIEW @this_model TO ROLE dev_role"],)defexecute(context,start,end,execution_time,**kwargs)->pd.DataFrame:returnpd.DataFrame([{"id":1,"name":"name"}])注意:这些语句的表名解析发生在虚拟层。在名为dev的环境中跑 plan 时,db.test_model和@this_model都会解析成db__dev.test_model,而不是物理表名。同样支持在model_defaults中配置项目级默认语句。
3. 蓝图(Blueprinting):一个模板批量生成多个模型
当多个模型逻辑相同、只是参数不同(比如每个客户一张表),不必复制粘贴 N 份代码——用blueprints属性传一个键值字典列表,一个文件就能"打印"出多个模型。
规则:模型名必须用蓝图里的变量做参数化,语法是@{变量名}。
importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,model@model("@{customer}.some_table",# 用蓝图变量 customer 参数化模型名kind="FULL",blueprints=[{"customer":"customer1","field_a":"x","field_b":"y"},{"customer":"customer2","field_a":"z","field_b":"w"},],columns={"field_a":"text","field_b":"text","customer":"text",},)defentrypoint(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:returnpd.DataFrame({"field_a":[context.blueprint_var("field_a")],"field_b":[context.blueprint_var("field_b")],"customer":[context.blueprint_var("customer")],})上面的定义会产生两个独立模型:customer1.some_table和customer2.some_table,各自使用对应的参数映射,变量通过context.blueprint_var读取。
3.1 动态生成蓝图列表
蓝图映射也可以由宏动态构造(适合从 CSV 等外部数据源读取清单):
@model("@{customer}.some_table",blueprints="@gen_blueprints()",...)fromsqlmeshimportmacro@macro()defgen_blueprints(evaluator):return("((customer := customer1, field_a := x, field_b := y),"" (customer := customer2, field_a := z, field_b := w))")还可以配合@EACH宏和全局列表变量(@values):
@model("@{customer}.some_table",blueprints="@EACH(@values, x -> (customer := schema_@x))",...)4. 模型属性里使用宏变量(小心 cron 的坑)
Python 模型的属性支持宏变量,但当宏变量出现在字符串内部时要特别小心。典型场景是把调度时间做成参数化 cron:
# 正确写法:整个表达式用引号包住,并加 @ 前缀@model("my_model",cron="@'*/@{mins} * * * *'",# 注意 @'...' 语法...)# 配合蓝图变量同样适用@model("@{customer}.scheduled_model",cron="@'0 @{hour} * * *'",blueprints=[{"customer":"customer_1","hour":2},# 凌晨 2 点跑{"customer":"customer_2","hour":8},# 早上 8 点跑],...)为什么要这么麻烦?因为 cron 表达式本身常用@表示别名(@daily、@hourly),会和 SQLMesh 的宏语法冲突,@'...'的写法能确保正确解析。
5. 全系列避坑清单(Best Practices)
| # | 规则 | 原因 |
|---|---|---|
| 1 | columns声明必须与实际返回的 DataFrame 完全一致 | SQLMesh 先建表再跑代码,schema 不符会引发意外行为 |
| 2 | 保持模型幂等 | 同一区间重跑结果必须一致,否则回刷数据时会产生脏数据 |
| 3 | 永远不要return空 DataFrame | 可能为空时改用条件yield(yield from ()) |
| 4 | 用 Spark/Snowflake/BigQuery 时优先返回对应原生 DataFrame | context.spark/snowpark/bigframe让计算分布式执行,避免本地内存瓶颈 |
| 5 | 读上游模型必先resolve_table | 直接硬编码表名会在 dev/prod 环境切换时拿错数据 |
| 6 | depends_on显式声明会覆盖函数体内的动态引用 | 避免依赖图与预期不符 |
| 7 | pre/post 语句避免创建物理表 | 多模型并发执行时会产生冲突 |
| 8 | 后置语句必须配合yield(不能return) | 后置语句要在函数产出数据之后执行 |
| 9 | 变量走函数参数时必须带默认值,且不藏在kwargs里 | 变量缺失时保证模型可加载 |
| 10 | 输出太大就用生成器分批yield | 降低单批内存占用 |
| 11 | 含宏变量的 cron 字符串用@'...'包裹 | 避免@符号解析冲突 |
| 12 | Python 模型不能用VIEW/SEED/MANAGED/EMBEDDEDkind | 需要这些 kind 时改用 SQL 模型 |
6. 汇总:一个串联全系列知识点的完整示例
把增量 kind、依赖解析、区间过滤、幂等产出写在一起,作为出师检验:
importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,modelfromsqlmesh.core.model.kindimportModelKindName@model("docs_example.final_demo",kind=dict(name=ModelKindName.INCREMENTAL_BY_TIME_RANGE,time_column="event_date",),columns={"id":"int","name":"text","event_date":"date",},depends_on=["docs_example.upstream_model"],cron="@daily",)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:# 1. 解析上游表名(自动登记依赖)table=context.resolve_table("docs_example.upstream_model")# 2. 只取本次时间区间内的数据 —— 保证增量与幂等df=context.fetchdf(f"SELECT id, name, event_date FROM{table}"f"WHERE event_date >= '{start}' AND event_date < '{end}'")# 3. 用 pandas 做业务逻辑df["name"]=df["name"].str.strip().str.lower()# 4. 可能为空的结果:用 yield 而不是 returnifdf.empty:yieldfrom()else:yielddf结语
三篇文章读完后,SQLMesh Python 模型的心法浓缩成三句话:
- 一个
@model装饰器 + 一个execute函数 = 一个模型,元数据字段与 SQL 模型一一对应; - schema 先于代码——
columns必填且必须与返回的 DataFrame 严格一致; - 让数据待在引擎里——能返回 Spark/Snowpark/Bigframe DataFrame 就不要落到 Pandas,输出太大就用生成器分批
yield。
你已经跨过了初学者到工程实践的门槛。下一步建议阅读官方文档的 model kinds 与 宏系统 章节,把增量策略和参数化能力用得更深。
参考资料:SQLMesh 官方文档 — Python models(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)