CertSafari
    Snowflake SnowPro Advanced: MLOps Engineer (MLA-B01)· Lessons

    Domain 4 · Lesson 12/17

    Building an ML pipeline with modular code and Snowflake ML Jobs

    Orchestrate end-to-end ML workflows.

    9 min read
    7.33% of exam
    2 sources
    Published 5 Oct 2026
    Docs as of 4 Oct 2026

    What you will be able to do

    • Restructure notebook code into parameterized stages that chain data validation, training, evaluation and registration
    • Submit pipeline steps as Snowflake ML Jobs from any Python environment using the snowflake-ml-python SDK
    • Choose between @remote, submit_file, submit_directory and submit_from_stage, and retrieve a job's result

    Key concept

    Task graph (DAG) as the pipeline backbone — An automated ML pipeline is a set of separate, parameterized stages (validate, train, evaluate, deploy) linked as a directed acyclic graph. The graph runs on a schedule or in response to an event, and each stage can hand heavy work off to dedicated compute.

    1.From notebook cells to a chained pipeline

    Snowflake splits an ML workflow into four broad stages: data exploration and preparation, data engineering, model development, and model deployment. Early work is usually iterative, done in a notebook or a local IDE. Once a model proves its value, the goal changes to operationalization. The pipeline gets hardened and automated, so that every change to code, data or model is built, tested and deployed the same way each time.

    The documentation recommends structuring things for this from the start. Parameterize inputs (tables, stages, hyperparameters) and give each step its own cell or function. Before you operationalize, turn each major step (data preparation, feature engineering, training, evaluation) into a function with clear inputs and outputs. Move every configuration value out of the code, so the same pipeline can run in dev and in prod. Also write one entrypoint script that runs the whole pipeline locally, which you can use for debugging.

    The body of the recommended run_pipeline.py entrypoint: a data validation gate runs before training, then evaluation, then registration, and every value comes from an environment-specific configpython
        # Load configuration
        config = load_config(args.config, args.env)
    
        # Execute pipeline stages
        raw_data = load_raw_data(config.data.source_table)
        validate_data_quality(raw_data, config.data.quality_checks)
        features = create_features(raw_data, config.features.transformations)
        model = train_model(features, config.model.hyperparameters)
        metrics = evaluate_model(model, features, config.model.eval_metrics)
        register_model(model, metrics, config.model.registry_name)

    This script is the chain the exam guide describes, written out in order. If validate_data_quality fails, training never starts. If a model has not been evaluated, it never reaches register_model. The --env argument (default dev) picks the configuration, so promoting the pipeline to production means changing an argument, not editing code. The example project layout puts this entrypoint in scripts/ next to a dag.py, with stage code under src/ml_pipeline/ split into data/, features/, models/ and inference/.

    The documentation names a Snowflake tool for each lifecycle stage. The rest of this lesson uses the ones in the table below.

    Snowflake tools for each lifecycle stage that an automated pipeline chains together
    Lifecycle stageToolRole in the pipeline
    Data EngineeringSnowpark DataFrames, UDFs/UDTFs, Feature StoreReproducible transforms at warehouse scale; consistent offline training sets
    Model TrainingML JobsOffload resource-intensive steps to high-memory, GPU or distributed compute
    Model DeploymentModel RegistryRegister and version models with lineage; safe promotion, audits and rollback
    Workflow OrchestrationScheduled NotebooksRun a parameterized notebook non-interactively on a schedule
    Workflow OrchestrationTask GraphsRun the pipeline as a DAG on a schedule or by event-based triggers

    Checkpoint 1 of 5· Check yourself

    A team wants the same training pipeline to run against dev tables during testing and prod tables after release, without editing source code. Which practice does Snowflake recommend?

    Sources1

    2.Running pipeline steps as ML Jobs from the Python SDK

    Training usually needs more compute than a warehouse provides. Snowflake ML Jobs run ML workloads inside Snowflake's container runtimes on compute pools, including GPU and high-memory CPU instances, and you can submit them from any development environment. The SDK is the snowflake-ml-python package, and ML Jobs need version 1.26.0 or later. The documentation also lists integration with orchestration tools such as Apache Airflow, so an external scheduler can call the same API.

    Most code written in Snowflake Notebooks works in ML Jobs without changes. One exception: some distributed ML APIs exist only inside the Container Runtime, so you must import them inside the job payload, not at the top of your local script.

    Every ML Jobs call needs a Snowpark Session. The documentation builds one with Session.builder.getOrCreate(), which reads a valid ~/.snowflake/config.toml. Once the session exists, you can pass it explicitly or let the API infer it from context.

    There are two ways to submit work. Function dispatch uses the @remote decorator. It serializes the function and its dependencies, uploads them to a stage, and runs them in a Container Runtime. File dispatch submits whole files or project directories, which suits multi-module projects, scripts that take command-line arguments, and existing projects that were not written for Snowflake. File dispatch has three methods:

    File dispatch: submit a single training script to a compute pool, with command-line argumentspython
    from snowflake.ml.jobs import submit_file
    
    # Run a single file
    job1 = submit_file(
      "train.py",
      "MY_COMPUTE_POOL",
      stage_name="payload_stage",
      args=["--data-table", "my_training_data"],
      session=session,
    )

    Checkpoint 2 of 5· Fill the gap

    This call submits a multi-file project directory. Which keyword tells the job which script to run?

    from snowflake.ml.jobs import submit_directory
    
    # Run from a directory
    job2 = submit_directory(
      "./ml_project/",
      "MY_COMPUTE_POOL",
       ? ="train.py",
      stage_name="payload_stage",
      session=session,
    )

    Checkpoint 3 of 5· Match them up

    Match each ML Jobs submission method to the situation it fits

    Tap a term, then the definition that fits it.

    Sources21

    3.Getting results back and pinning the runtime

    A pipeline can only chain steps if each step's output is available to the next. Every submission method returns an MLJob object that you use to manage and monitor the run. Inside the job, a Snowpark Session is already available, so the payload can read tables and stages and write results back. For function dispatch, you can declare a session: Session parameter and the session is injected automatically. That only works for function payloads. Session.builder.getOrCreate() works for every payload type.

    How you return a value depends on the dispatch style. A @remote function just returns it. A file-based job assigns it to the special __return__ variable, for example __return__ = main(). On the client, MLJob.result() blocks until the job reaches a terminal state. It then returns the payload's value, or raises an exception if the run failed. If the payload returned nothing, a successful job gives None.

    Pinning a function-dispatch job to Container Runtime 2.3.0, so an automated pipeline does not pick up a newer runtime unexpectedlypython
    from snowflake.ml.jobs import remote
    
    @remote("MY_COMPUTE_POOL", stage_name="payload_stage", session=session, runtime_environment="2.3.0")
    def train_model(data_table: str):
      # Provide your ML code here, including imports and function calls
      ...

    If you leave out runtime_environment, the job uses the latest Container Runtime available on the compute pool. That is fine for experiments, but a production pipeline should pin a full version string so runs are reproducible. For package or compliance needs, you can register a custom image and reference it as runtime_environment="cre@<name>". Compute pools are sized with MIN_NODES and MAX_NODES. The default pool uses the CPU_X64_S instance family, with a minimum of 1 node and a maximum of 25.

    Checkpoint 4 of 5· Check yourself

    A file-based ML Job runs train.py, which sets __return__ = main(). The orchestrating script calls job.result() while the job is still running. What does result() do?

    Checkpoint 5 of 5· Exam question

    A nightly task graph chains VALIDATE_DATA, TRAIN_MODEL, and DEPLOY_MODEL. When the validation step finds a null rate above the agreed threshold, training and deployment must not run, and the failure must show up in task history. How should the validation step be built?

    Sources2

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 1.ML Jobs can only be launched from a Snowflake Notebook or worksheet.Why is that wrong?

      You can submit ML Jobs from any development environment with snowflake-ml-python, including local IDEs and external orchestrators such as Airflow.

      Covered in Running pipeline steps as ML Jobs from the Python SDK

    2. 2.A file-based ML Job returns whatever its main() function returns, the same way a @remote function does.Why is that wrong?

      A file payload has to assign its result to the special __return__ variable. Otherwise result() gives None on success.

      Covered in Getting results back and pinning the runtime

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 1.
      “Parameterize inputs (tables, stages, hyperparameters) and keep steps modular for portability.”
      ↩︎ From notebook cells to a chained pipeline
      “We recommend also authoring an entrypoint script which executes the end-to-end pipeline locally for debugging and future development.”
      ↩︎ From notebook cells to a chained pipeline
      “These APIs are available inside ML Jobs, but need to be imported inside the ML Job payload.”
      ↩︎ Running pipeline steps as ML Jobs from the Python SDK
      “Operationalize your ML pipeline into a Directed Acyclic Graph (DAG) and configure it to run on a schedule or by event based triggers.”
      ↩︎ Key concept
      “Parameterize all configuration values like table names and hyperparameters to enable cross-environment deployment.”
      ↩︎ Checkpoint
    2. 2.
      “Snowflake ML Jobs are available in snowflake-ml-python version 1.26.0 and later.”
      ↩︎ Running pipeline steps as ML Jobs from the Python SDK
      “Integrate with orchestration tools, such as Apache Airflow.”
      ↩︎ Running pipeline steps as ML Jobs from the Python SDK
      “When running ML Jobs on Snowflake, a Snowpark Session is automatically available in the execution context.”
      ↩︎ Getting results back and pinning the runtime
      “Snowflake automatically uses the latest available version of the Snowflake Container Runtime on your compute pool.”
      ↩︎ Getting results back and pinning the runtime
      “You can run them from any development environment. You don’t need to run the code in a Snowflake worksheet or notebook.”
      ↩︎ Exam trap 1
      “For file-based jobs, use the special __return__ variable to specify the return value.”
      ↩︎ Exam trap 2
      “ls = list_jobs() # This will fail! You must create a session first.”
      ↩︎ Prediction
      “submit_from_stage: For running Python projects saved on a Snowflake stage”
      ↩︎ Checkpoint
      “The API blocks the calling thread until the job reaches a terminal state, then returns the payload’s return value”
      ↩︎ Checkpoint

    Continue to page 2 of 2

    Orchestrating ML pipelines with task graphs and Snowflake CLI

    Spotted a mistake, or was something unclear? Tell us.