:依赖管理、引擎实战与内存优化)
本系列基于 SQLMesh 官方文档https://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/整理共 3 篇面向初学者。本篇是第二篇解决数据怎么取、算在哪、内存爆了怎么办三个实战问题。第一篇基础语法与核心概念第三篇前后置语句、蓝图批量建模、避坑清单1. 依赖管理resolve_table与depends_on想读取上游模型的数据必须先解析出它在当前环境下的真实表名这就是resolve_table的职责tablecontext.resolve_table(docs_example.upstream_model)dfcontext.fetchdf(fSELECT * FROM{table})resolve_table有两重作用返回当前运行环境如 dev 环境的schema__dev前缀下正确的表名自动把被引用的模型登记为当前模型的依赖。另一种声明依赖的方式是在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,kindFULL,columns{id:INT},)defexecute(context,**kwargs):tablecontext.resolve_table(f{context.var(schema_name)}.test_model1)select_queryexp.select(*).from_(table)returncontext.fetchdf(select_query)2. 实战示例三连Basic → SQLPandas → PySpark2.1 查询上游模型 Pandas 处理最常用的一种模式SQL 负责取数和粗筛pandas 负责灵活加工importtypingastfromdatetimeimportdatetimeimportpandasaspdfromsqlmeshimportExecutionContext,modelmodel(docs_example.sql_pandas,columns{id:int,name:text,},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)-pd.DataFrame:# 获取上游模型的表名并自动登记为依赖tablecontext.resolve_table(upstream_model)# 把数据取回来如果引擎是 Spark这里返回的就是 Spark DataFramedfcontext.fetchdf(fSELECT id, name FROM{table})# 做一些 pandas 擅长的事df[id]1returndf2.2 PySpark 示例分布式计算的正确姿势如果你使用 Spark 引擎推荐直接用 Spark DataFrame API 而不是 Pandas——数据全程在集群分布式计算不会拉到本地。importtypingastfromdatetimeimportdatetimefrompyspark.sqlimportDataFrame,functionsfromsqlmeshimportExecutionContext,modelmodel(docs_example.pyspark,columns{id:int,name:text,country:text,},)defexecute(context:ExecutionContext,start:datetime,end:datetime,execution_time:datetime,**kwargs:t.Any,)-DataFrame:# 获取上游模型表名并登记依赖tablecontext.resolve_table(upstream_model)# 用 Spark DataFrame API 增加一列 countrydfcontext.spark.table(table).withColumn(country,functions.lit(USA))# 直接返回 PySpark DataFrame本地不计算任何数据returndf三个关键学习点通过context.spark拿到 SparkSession再用spark.table(表名)载入数据比先fetchdf转成 Pandas 高效得多返回类型注解写pyspark.sql.DataFrame而不是pd.DataFrame只要最终返回的是 Spark DataFrame计算就发生在集群上没有本地内存瓶颈。3. 同样思路的另外两种引擎Snowpark 与 Bigframe3.1 SnowparkSnowflake 引擎用context.snowpark操作 DataFrame计算下推到 Snowflakeimporttypingastfromdatetimeimportdatetimefromsnowflake.snowpark.dataframeimportDataFramefromsqlmeshimportExecutionContext,modelmodel(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本地不计算任何数据dfcontext.snowpark.create_dataframe([[1,a,usa],[2,b,cad]],schema[id,name,country])dfdf.filter(df.id1)returndf3.2 BigframeBigQuery 引擎用context.bigframe所有计算都在 BigQuery 完成。它甚至支持把本地 Python 函数注册为 remote function 在集群上执行importtypingastfromdatetimeimportdatetimefrombigframes.pandasimportDataFramefromsqlmeshimportExecutionContext,modeldefget_bucket(num:int):ifnotnum:returnNAboundary10returnat_or_above_10ifnumboundaryelsebelow_10model(mart.wiki,columns{title:text,views:int,bucket:text,},)defexecute(context,start,end,execution_time,**kwargs)-DataFrame:# 把本地 Python 函数包装成 BigQuery remote functionremote_get_bucketcontext.bigframe.remote_function([int],str)(get_bucket)# 只返回 Bigframe 句柄数据不落到本地dfcontext.bigframe.read_gbq(bigquery-samples.wikipedia_pageviews.200809h)df(df[df.title.str.contains(r[Gg]oogle)].groupby([title],as_indexFalse)[views].sum(numeric_onlyTrue).sort_values(views,ascendingFalse))returndf.assign(bucketdf[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:tablecontext.resolve_table(upstream_model)foriinrange(3):# 分 3 次查询每次只取一块数据避免内存耗尽dfcontext.fetchdf(fSELECT id from{table}WHERE id {i})yielddf5. 空表禁忌绝对不能return空 DataFramePython 模型不允许返回空的 DataFrame。如果你的代码有可能产出空结果必须改成条件yieldmodel(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重计算会下推到集群本地只做编排。小结读上游模型必先resolve_table——硬编码表名会在 dev/prod 切换时拿错数据depends_on显式声明优先于函数体内的动态引用引擎有原生 DataFrame APIspark/snowpark/bigframe时尽量让计算留在集群别拉回本地输出太大用生成器分批yield可能为空的结果绝不用return空表。下一篇讲工程化细节前后置语句、宏变量在属性中的坑、以及用蓝图一次生成多个模型。参考资料SQLMesh 官方文档 — Python modelshttps://sqlmesh.readthedocs.io/en/stable/concepts/models/python_models/