ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SQLMesh 实战:SQL 模型与 Python 模型混合使用指南

SQLMesh 实战:SQL 模型与 Python 模型混合使用指南 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.yamlgateways:local:connection:type:duckdbdatabase:db.dbdefault_gateway:localmodel_defaults:dialect:duckdbSQL 模型加工订单、库存、物料订单、库存、物料主数据适合直接用 SQL 模型。它们结构清晰主要工作是字段标准化、聚合和关联。models/fct_order.sqlMODEL(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.sqlMODEL(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.sqlMODEL(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.sqlMODEL(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_hand0THENout_of_stockWHENCOALESCE(o.confirmed_qty,0)i.qty_on_handTHENshortageWHENi.qty_on_hand150THENover_stockELSEnormalENDASstock_statusFROMdemo.fct_inventoryASiLEFTJOINdemo.dim_materialASmONi.material_codem.material_codeLEFTJOIN(SELECTmaterial_code,SUM(CASEWHENorder_statusconfirmedTHENorder_qtyELSE0END)ASconfirmed_qty,SUM(CASEWHENorder_statuspendingTHENorder_qtyELSE0END)ASpending_qtyFROMdemo.fct_orderGROUPBYmaterial_code)ASoONi.material_codeo.material_code;这类逻辑仍然适合 SQL。因为它本质是多表关联和条件判断用 SQL 表达更自然也方便后续排查。Python 模型计算库存健康评分当规则变得复杂时就可以转入 Python 模型。例如库存健康评分不仅要看库存数量还要结合物料状态、订单覆盖倍数、是否缺货、是否超储等因素。models/inventory_health_score.pyimportpandasaspdfromsqlmeshimportmodelmodel(demo.inventory_health_score,kindFULL,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):dfcontext.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:coveragerow[qty_on_hand]/row[confirmed_qty]ifcoverage3:return90ifcoverage1:return75return60return50df[health_score]df.apply(score,axis1)returndfPython 模型可以从上游宽表读取数据按业务规则计算评分再输出结果表。后续如果评分规则需要调整也只需修改 Python 代码不影响 SQL 层。Python 模型接入 Great Expectations 质量门禁Great Expectations 不建议写进 SQL 文件。更合适的方式是放在依赖链最后的 Python 模型或者外部调度脚本中。models/ge_quality_gate.pyimportpandasaspdimportgreat_expectationsasgxfromsqlmeshimportmodelmodel(demo.ge_quality_gate,kindFULL,depends_on{demo.fct_order,demo.fct_inventory,demo.dim_material,demo.inventory_health_score,},)defexecute(context,**kwargs):ctxgx.get_context()dsctx.sources.add_or_update_sqlalchemy(nameduckdb_demo,connection_stringduckdb:///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_passedTrueforiteminchecks:tableitem[table]assetds.add_table_asset(nametable,table_nametable)batchasset.build_batch_request()validatorctx.get_validator(batch_requestbatch,expectation_suite_namef{table}_suite,)formethod_name,kwinitem[rules]:getattr(validator,method_name)(**kw)resultvalidator.validate()ifnotresult.success:all_passedFalsefailed[r.expectation_config.expectation_typeforrinresult.resultsifnotr.success]raiseRuntimeError(fGE 校验失败{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 完成评分和校验最后形成稳定、可观测的数据管道。
返回列表