SQLMesh 同时支持 SQL 模型和 Python 模型。实际项目中,订单、库存、物料主数据这类结构化加工应优先使用 SQL 模型;评分规则、外部接口、机器学习、数据质量门禁等复杂逻辑再交给 Python 模型。本文通过一个完整可跑示例,展示两者如何配合,并给出生产环境落地建议。
为什么需要混合使用
数据管道里大部分工作仍是表关联、聚合、清洗和指标计算。这些场景用 SQL 模型更直观,也更容易被业务、实施和数据分析人员理解。
Python 模型的价值在于处理 SQL 不擅长的事。例如调用外部价格接口、执行复杂评分规则、运行机器学习模型、接入 Great Expectations 做质量门禁,或者把校验结果写入外部系统。
混合使用的好处是分工清楚。SQL 负责数据加工,Python 负责复杂逻辑和外部联动。这样既保持模型可读,又不会把复杂代码硬塞进 SQL 里。
案例背景
以制造业常见的订单、库存、物料主数据为例。目标是先完成基础数据加工,再计算库存健康评分,最后通过 Great Expectations 进行质量校验。
环境准备
pipinstallsqlmesh great-expectations duckdb pandas项目结构:
sqlmesh-ge-demo/ ├── config.yaml ├── models/ │ ├── dim_material.sql │ ├── fct_order.sql │ ├── fct_inventory.sql │ ├── inventory_health_base.sql │ ├── inventory_health_score.py │ └── ge_quality_gate.py └── run_pipeline.shconfig.yaml:
gateways:local:connection:type:duckdbdatabase:db.dbdefault_gateway:localmodel_defaults:dialect:duckdbSQL 模型:加工订单、库存、物料
订单、库存、物料主数据适合直接用 SQL 模型。它们结构清晰,主要工作是字段标准化、聚合和关联。
models/fct_order.sql:
MODEL(name demo.fct_order,kindFULL,grain order_id,audits(NOT_NULL(columns=(order_id,customer_code,material_code,order_qty,amount)),UNIQUE_VALUES(columns=(order_id)),accepted_values(order_status,['pending','confirmed','shipped','cancelled'])));SELECT*FROM(VALUES('PO202610010001','CUST001','M001',10.00,42000.00,'confirmed','2026-10-01'),('PO202610010002','CUST002','M002',2.50,39500.00,'confirmed','2026-10-01'),('PO202610010003','CUST001','M003',5.00,28000.00,'pending','2026-10-01'),('PO202610010004','CUST003','M001',20.00,84000.00,'shipped','2026-09-30'),('PO202610010005','CUST004','M005',3.00,58500.00,'confirmed','2026-10-01'),('PO202610010006','CUST002','M003',8.00,44800.00,'cancelled','2026-09-29'))ASt(order_id,customer_code,material_code,order_qty,amount,order_status,order_date);models/fct_inventory.sql:
MODEL(name demo.fct_inventory,kindFULL,grain(warehouse_code,material_code),audits(NOT_NULL(columns=(warehouse_code,material_code,qty_on_hand)),UNIQUE_VALUES(columns=(warehouse_code,material_code))));SELECT*FROM(VALUES('WH01','M001',100.00),('WH01','M002',20.00),('WH01','M003',50.00),('WH02','M001',200.00),('WH02','M005',30.00))ASt(warehouse_code,material_code,qty_on_hand);models/dim_material.sql:
MODEL(name demo.dim_material,kindFULL,grain material_code,audits(NOT_NULL(columns=(material_code,material_name,category,unit_price,status)),UNIQUE_VALUES(columns=(material_code)),accepted_values(category,['钢材','有色金属']),accepted_values(status,['active','inactive'])));SELECT*FROM(VALUES('M001','圆钢HRB400','钢材',4200.00,'active'),('M002','不锈钢板304','钢材',15800.00,'active'),('M003','镀锌钢管','钢材',5600.00,'active'),('M004','铜管T2','有色金属',68000.00,'inactive'),('M005','铝卷1060','有色金属',19500.00,'active'))ASt(material_code,material_name,category,unit_price,status);SQL 模型:构建库存健康宽表
在基础表完成后,可以用 SQL 继续构建库存健康宽表。它会关联物料主数据、库存数据和订单数据,计算已确认订单量、待确认订单量,并判断库存状态。
models/inventory_health_base.sql:
MODEL(name demo.inventory_health_base,kindFULL,grain(warehouse_code,material_code));SELECTi.warehouse_code,i.material_code,m.material_name,m.category,m.unit_price,m.statusASmaterial_status,i.qty_on_hand,COALESCE(o.confirmed_qty,0)ASconfirmed_qty,COALESCE(o.pending_qty,0)ASpending_qty,CASEWHENi.qty_on_hand=0THEN'out_of_stock'WHENCOALESCE(o.confirmed_qty,0)>i.qty_on_handTHEN'shortage'WHENi.qty_on_hand>150THEN'over_stock'ELSE'normal'ENDASstock_statusFROMdemo.fct_inventoryASiLEFTJOINdemo.dim_materialASmONi.material_code=m.material_codeLEFTJOIN(SELECTmaterial_code,SUM(CASEWHENorder_status='confirmed'THENorder_qtyELSE0END)ASconfirmed_qty,SUM(CASEWHENorder_status='pending'THENorder_qtyELSE0END)ASpending_qtyFROMdemo.fct_orderGROUPBYmaterial_code)ASoONi.material_code=o.material_code;这类逻辑仍然适合 SQL。因为它本质是多表关联和条件判断,用 SQL 表达更自然,也方便后续排查。
Python 模型:计算库存健康评分
当规则变得复杂时,就可以转入 Python 模型。例如库存健康评分不仅要看库存数量,还要结合物料状态、订单覆盖倍数、是否缺货、是否超储等因素。
models/inventory_health_score.py:
importpandasaspdfromsqlmeshimportmodel@model("demo.inventory_health_score",kind="FULL",depends_on={"demo.inventory_health_base"},columns={"warehouse_code":"TEXT","material_code":"TEXT","material_name":"TEXT","category":"TEXT","unit_price":"DOUBLE","material_status":"TEXT","qty_on_hand":"DOUBLE","confirmed_qty":"DOUBLE","pending_qty":"DOUBLE","stock_status":"TEXT","health_score":"DOUBLE",},)defexecute(context,**kwargs):df=context.fetchdf("SELECT * FROM demo.inventory_health_base")defscore(row):ifrow["material_status"]!="active":return0ifrow["stock_status"]=="out_of_stock":return20ifrow["stock_status"]=="shortage":return40ifrow["stock_status"]=="over_stock":return60ifrow["qty_on_hand"]>0androw["confirmed_qty"]>0:coverage=row["qty_on_hand"]/row["confirmed_qty"]ifcoverage>=3:return90ifcoverage>=1:return75return60return50df["health_score"]=df.apply(score,axis=1)returndfPython 模型可以从上游宽表读取数据,按业务规则计算评分,再输出结果表。后续如果评分规则需要调整,也只需修改 Python 代码,不影响 SQL 层。
Python 模型:接入 Great Expectations 质量门禁
Great Expectations 不建议写进 SQL 文件。更合适的方式是放在依赖链最后的 Python 模型,或者外部调度脚本中。
models/ge_quality_gate.py:
importpandasaspdimportgreat_expectationsasgxfromsqlmeshimportmodel@model("demo.ge_quality_gate",kind="FULL",depends_on={"demo.fct_order","demo.fct_inventory","demo.dim_material","demo.inventory_health_score",},)defexecute(context,**kwargs):ctx=gx.get_context()ds=ctx.sources.add_or_update_sqlalchemy(name="duckdb_demo",connection_string="duckdb:///db.db",)checks=[{"table":"demo.fct_order","rules":[("expect_column_values_to_not_be_null",{"column":"order_id"}),("expect_column_values_to_be_unique",{"column":"order_id"}),("expect_column_values_to_be_greater_than_or_equal_to",{"column":"amount","value":0}),],},{"table":"demo.fct_inventory","rules":[("expect_column_values_to_not_be_null",{"column":"warehouse_code"}),("expect_column_values_to_not_be_null",{"column":"material_code"}),("expect_column_values_to_be_greater_than_or_equal_to",{"column":"qty_on_hand","value":0}),],},{"table":"demo.inventory_health_score","rules":[("expect_column_values_to_not_be_null",{"column":"material_code"}),("expect_column_values_to_be_between",{"column":"health_score","min_value":0,"max_value":100}),],},]all_passed=Trueforiteminchecks:table=item["table"]asset=ds.add_table_asset(name=table,table_name=table)batch=asset.build_batch_request()validator=ctx.get_validator(batch_request=batch,expectation_suite_name=f"{table}_suite",)formethod_name,kwinitem["rules"]:getattr(validator,method_name)(**kw)result=validator.validate()ifnotresult.success:all_passed=Falsefailed=[r.expectation_config.expectation_typeforrinresult.resultsifnotr.success]raiseRuntimeError(f"GE 校验失败{table}:{failed}")ctx.build_data_docs()returnpd.DataFrame({"status":["ge_passed"]})如果 GE 校验失败,模型会抛异常,SQLMesh 会阻断后续依赖,下游任务就不会继续消费错误数据。
执行方式
sqlmesh plan --auto-applySQLMesh 会按依赖执行:
fct_order / fct_inventory / dim_material → inventory_health_base → inventory_health_score → ge_quality_gate如果希望用外部脚本串联,也可以这样:
run_pipeline.sh:
#!/usr/bin/env bashset-euopipefailecho"=== 1. SQLMesh 物化模型 ==="sqlmesh plan --auto-applyecho"=== 2. 运行 Great Expectations 校验 ==="python ge_suite/run_ge_checks.pyecho"=== 3. 完成 ==="生产环境配置
如果生产环境使用 PostgreSQL、Trino 或 Snowflake,不需要改变模型结构。只需调整 SQLMesh 网关配置和 Great Expectations 的连接字符串即可。
SQL 模型仍负责主要数据加工,Python 模型负责复杂逻辑和质量门禁。两者通过依赖关系串联,执行顺序清晰,问题也容易定位。
总结
SQL 模型适合处理结构化数据加工,Python 模型适合处理复杂规则、外部接口、机器学习和质量门禁。两者混合使用时,应让 SQL 承担主要计算,Python 承担增强能力。
这种模式的优势是结构清晰、可维护性强、易于扩展。订单、库存、物料主数据等场景可以先用 SQL 打好基础,再用 Python 完成评分和校验,最后形成稳定、可观测的数据管道。