news 2026/10/3 7:16:47

SQLMesh Python 模型入门(二):依赖管理、引擎实战与内存优化

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SQLMesh Python 模型入门(二):依赖管理、引擎实战与内存优化

本系列基于 SQLMesh 官方文档(https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/)整理,共 3 篇,面向初学者。本篇是第二篇,解决"数据怎么取、算在哪、内存爆了怎么办"三个实战问题。

  • 第(一)篇:基础语法与核心概念
  • 第(三)篇:前后置语句、蓝图批量建模、避坑清单

1. 依赖管理:resolve_table与depends_on

想读取上游模型的数据,必须先解析出它在当前环境下的真实表名,这就是resolve_table的职责:

table=context.resolve_table("docs_example.upstream_model")df=context.fetchdf(f"SELECT * FROM{table}")

resolve_table有两重作用:

  1. 返回当前运行环境(如 dev 环境的schema__dev前缀)下正确的表名;
  2. 自动把被引用的模型登记为当前模型的依赖。

另一种声明依赖的方式是在@model装饰器里显式写depends_on。规则是:装饰器里显式声明的依赖优先于函数体内的动态引用。看这个官方例子:

@model("my_model.with_explicit_dependencies",depends_on=["docs_example.upstream_dependency"],# ✅ 会被捕获)defexecute(context,start,end,execution_time,**kwargs):# ❌ 由于装饰器里已声明依赖,这里的引用会被忽略context.resolve_table("docs_example.another_dependency")...

此外,用户自定义的全局变量或蓝图变量也可以出现在resolve_table的调用中:

@model("@schema_name.test_model2",kind="FULL",columns={"id":"INT"},)defexecute(context,**kwargs):table=context.resolve_table(f"{context.var('schema_name')}.test_model1")select_query=exp.select("*").from_(table)returncontext.fetchdf(select_query)

2. 实战示例三连:Basic → SQL+Pandas → PySpark

2.1 查询上游模型 + Pandas 处理

最常用的一种模式:SQL 负责取数和粗筛,pandas 负责灵活加工:

importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,model@model("docs_example.sql_pandas",columns={"id":"int","name":"text",},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:# 获取上游模型的表名,并自动登记为依赖table=context.resolve_table("upstream_model")# 把数据取回来;如果引擎是 Spark,这里返回的就是 Spark DataFramedf=context.fetchdf(f"SELECT id, name FROM{table}")# 做一些 pandas 擅长的事df["id"]+=1returndf

2.2 PySpark 示例:分布式计算的正确姿势

如果你使用 Spark 引擎,推荐直接用 Spark DataFrame API 而不是 Pandas——数据全程在集群分布式计算,不会拉到本地。

importtypingastfromdatetimeimportdatetimefrompyspark.sqlimportDataFrame,functionsfromsqlmeshimportExecutionContext,model@model("docs_example.pyspark",columns={"id":"int","name":"text","country":"text",},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->DataFrame:# 获取上游模型表名并登记依赖table=context.resolve_table("upstream_model")# 用 Spark DataFrame API 增加一列 countrydf=context.spark.table(table).withColumn("country",functions.lit("USA"))# 直接返回 PySpark DataFrame,本地不计算任何数据returndf

三个关键学习点:

  1. 通过context.spark拿到 SparkSession,再用spark.table(表名)载入数据,比先fetchdf转成 Pandas 高效得多;
  2. 返回类型注解写pyspark.sql.DataFrame而不是pd.DataFrame;
  3. 只要最终返回的是 Spark DataFrame,计算就发生在集群上,没有本地内存瓶颈。

3. 同样思路的另外两种引擎:Snowpark 与 Bigframe

3.1 Snowpark(Snowflake 引擎)

用context.snowpark操作 DataFrame,计算下推到 Snowflake:

importtypingastfromdatetimeimportdatetimefromsnowflake.snowpark.dataframeimportDataFramefromsqlmeshimportExecutionContext,model@model("docs_example.snowpark",columns={"id":"int","name":"text","country":"text",},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->DataFrame:# 直接返回 snowpark DataFrame,本地不计算任何数据df=context.snowpark.create_dataframe([[1,"a","usa"],[2,"b","cad"]],schema=["id","name","country"])df=df.filter(df.id>1)returndf

3.2 Bigframe(BigQuery 引擎)

用context.bigframe,所有计算都在 BigQuery 完成。它甚至支持把本地 Python 函数注册为 remote function 在集群上执行:

importtypingastfromdatetimeimportdatetimefrombigframes.pandasimportDataFramefromsqlmeshimportExecutionContext,modeldefget_bucket(num:int):ifnotnum:return"NA"boundary=10return"at_or_above_10"ifnum>=boundaryelse"below_10"@model("mart.wiki",columns={"title":"text","views":"int","bucket":"text",},)defexecute(context,start,end,execution_time,**kwargs)->DataFrame:# 把本地 Python 函数包装成 BigQuery remote functionremote_get_bucket=context.bigframe.remote_function([int],str)(get_bucket)# 只返回 Bigframe 句柄,数据不落到本地df=context.bigframe.read_gbq("bigquery-samples.wikipedia_pageviews.200809h")df=(df[df.title.str.contains(r"[Gg]oogle")].groupby(["title"],as_index=False)["views"].sum(numeric_only=True).sort_values("views",ascending=False))returndf.assign(bucket=df["views"].apply(remote_get_bucket))

四种返回类型选择速查:

场景返回类型入口计算位置
通用 / 小数据量Pandas DataFramecontext.fetchdf本地内存
Spark 引擎PySpark DataFramecontext.spark集群分布式
Snowflake 引擎Snowpark DataFramecontext.snowparkSnowflake 内
BigQuery 引擎Bigframe DataFramecontext.bigframeBigQuery 内

4. 大输出分块:用生成器yield降低内存占用

Pandas 是单机内存型框架,数据量太大时内存会爆,又用不了 Spark 时,SQLMesh 允许用 Python 生成器yield把输出拆成多批,每次只把一小块数据加载进内存:

@model("docs_example.batching",columns={"id":"int",},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)->pd.DataFrame:table=context.resolve_table("upstream_model")foriinrange(3):# 分 3 次查询,每次只取一块数据,避免内存耗尽df=context.fetchdf(f"SELECT id from{table}WHERE id ={i}")yielddf

5. 空表禁忌:绝对不能return空 DataFrame

Python 模型不允许返回空的 DataFrame。如果你的代码有可能产出空结果,必须改成条件yield:

@model("my_model.empty_df")defexecute(context:ExecutionContext,)->pd.DataFrame:# ... 生成 df 的代码 ...ifdf.empty:yieldfrom()# 空的话什么都不产出else:yielddf

记住口诀:有数据 →yield df;可能没数据 → 永远不要return空表。


6. 序列化:代码到底在哪里运行?

SQLMesh 通过自研的序列化框架,在运行 SQLMesh 的机器本地执行 Python 代码。这意味着:

  • 你的 Python 环境(依赖包版本)需要就绪;
  • 如果引擎是 Spark/Snowflake/BigQuery 且你返回的是对应 DataFrame,重计算会下推到集群,本地只做编排。

小结

  1. 读上游模型必先resolve_table——硬编码表名会在 dev/prod 切换时拿错数据;
  2. depends_on显式声明优先于函数体内的动态引用;
  3. 引擎有原生 DataFrame API(spark/snowpark/bigframe)时,尽量让计算留在集群,别拉回本地;
  4. 输出太大用生成器分批yield,可能为空的结果绝不用return空表。

下一篇讲工程化细节:前后置语句、宏变量在属性中的坑、以及用蓝图一次生成多个模型。


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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/3 7:15:41

图像分割目标检测

1. 视觉识别任务概述视觉识别是计算机视觉领域的核心任务,旨在让计算机理解图像内容。根据输出粒度的不同,视觉识别任务可以分为四大类:图像分类、语义分割、目标检测和实例分割。这四类任务对图像的理解程度逐层加深,从「整张图是…

作者头像 李华
网站建设 2026/10/3 7:15:41

GCDBZ 高程点布置插件 V3.0:一个命令把 CASS 高程点批量排出来

其实「已知一点高程 设计坡度,沿方向布置一串高程点」这件事,完全可以用一个插件自动完成。今天这篇,就把 GCDBZ 高程点布置插件 V3.0 讲透:它解决什么问题、按什么原理工作、怎么用、有哪些坑要避开。01 它解决什么问题&#xf…

作者头像 李华
网站建设 2026/10/3 7:15:11

多相机采集系统搭建:从传感器选型到存储落盘的工程链路

多相机采集系统(EGO 数据采集、多目同步采集等)的搭建,本质上是把一条从传感器到存储的数据链路做通。它的难点很少落在"相机能不能出图"上,而集中在链路中后段:接口带宽是否够、多路汇聚会不会互相抢、时间…

作者头像 李华
网站建设 2026/10/3 7:14:04

微信小程序+SSM驾考系统开发全攻略:从毕设开题到答辩排坑

最近被问得最多的一个毕设题目,就是“基于微信小程序的优选驾考SSM”。微信小程序做前端,SSM(Spring SpringMVC MyBatis)做后端,搭一套驾考学习、刷题、预约考试的系统。这个组合在近几年的毕业设计里出现频率非常高…

作者头像 李华