# 02 - Vertex Pipelines with R  



#### Install Vertex SDK for Python

We will be using [Vertex SDK for Python](https://cloud.google.com/vertex-ai/docs/start/client-libraries#python) to interact with Vertex AI services. The high-level aiplatform library is designed to simplify common data science workflows by using wrapper classes and opinionated defaults.

In [None]:
!pip -q install --user --upgrade kfp
!pip -q install --user --upgrade google-cloud-pipeline-components 
!pip -q install --user --upgrade google-cloud-aiplatform

### Restart the kernel

After you install the additional packages, you need to restart the notebook kernel so it can find the packages.

In [None]:
# Automatically restart kernel after installs
import os

if not os.getenv("IS_TESTING"):
    # Automatically restart kernel after installs
    import IPython

    app = IPython.Application.instance()
    app.kernel.do_shutdown(True)

Check the versions of the packages you installed.  The KFP SDK version should be >=1.6.

In [None]:
! python3 -c "import kfp; print('kfp version: {}'.format(kfp.__version__))"
! python3 -c "import google_cloud_pipeline_components; print('google_cloud_pipeline_components version: {}'.format(google_cloud_pipeline_components.__version__))"

## Before you begin

#### Set your project ID

**If you don't know your project ID**, you may be able to get your project ID using `gcloud`.

In [1]:
import os
# Get your Google Cloud project ID from gcloud
shell_output=!gcloud config list --format 'value(core.project)' 2>/dev/null

try:
    PROJECT_ID = shell_output[0]
except IndexError:
    PROJECT_ID = None

# Get your Google Cloud project ID from gcloud
if not os.getenv("IS_TESTING"):
    shell_output=!gcloud config list --format 'value(core.project)' 2>/dev/null
    PROJECT_ID = shell_output[0]
    print("Project ID: ", PROJECT_ID)

Project ID:  demos-vertex-ai


#### Timestamp

If you are in a live tutorial session, you might be using a shared test account or project. To avoid name collisions between users on resources created, you create a timestamp for each instance session, and append it onto the name of resources you create in this tutorial.

In [2]:
from datetime import datetime

def get_timestamp():
    return datetime.now().strftime("%Y%m%d%H%M%S")

TIMESTAMP = get_timestamp()
print(f"TIMESTAMP = {TIMESTAMP}")

TIMESTAMP = 20221110202015


### Create a Cloud Storage bucket as necessary



In [3]:
BUCKET_NAME = "gs://[your-bucket-name]"  # @param {type:"string"}
REGION = "us-central1"  # @param {type:"string"}

In [4]:
if BUCKET_NAME == "" or BUCKET_NAME is None or BUCKET_NAME == "gs://[your-bucket-name]":
    BUCKET_NAME = "gs://" + PROJECT_ID + "-aip-r-model-02-" + TIMESTAMP
BUCKET_NAME

'gs://demos-vertex-ai-aip-r-model-02-20221110202015'

**Only if your bucket doesn't already exist**: Run the following cell to create your Cloud Storage bucket.

In [5]:
! gsutil mb -l $REGION $BUCKET_NAME

Creating gs://demos-vertex-ai-aip-r-model-02-20221110202015/...


Finally, validate access to your Cloud Storage bucket by examining its contents:

In [6]:
! gsutil ls -al $BUCKET_NAME

### Import libraries and define constants

Define some constants. See the "Before you begin" section of the Managed Pipelines User Guide for information on creating your API key.

In [7]:
APP_NAME = "r-model-02"

Do some imports:

In [9]:
import json
from typing import NamedTuple, List

from google_cloud_pipeline_components import aiplatform as aip_components
from google_cloud_pipeline_components.v1.custom_job import CustomTrainingJobOp
        
from google_cloud_pipeline_components.types import artifact_types
from google.cloud import aiplatform
from google.cloud.aiplatform import pipeline_jobs

from kfp.v2 import compiler
from kfp.v2 import dsl
from kfp.v2 import components

In [10]:
! mkdir ./src

mkdir: cannot create directory ‘./src’: File exists


#### Enable Artifact Registry API

First, you must enable the Artifact Registry API service for your project.

Learn more about [Enabling service](https://cloud.google.com/artifact-registry/docs/enable-service).

In [11]:
! gcloud services enable artifactregistry.googleapis.com

#### Create a private Docker repository

Your first step is to create your own Docker repository in Google Artifact Registry.

1. Run the `gcloud artifacts repositories create` command to create a new Docker repository with your region with the description "docker repository".

2. Run the `gcloud artifacts repositories list` command to verify that your repository was created.

In [12]:
PRIVATE_REPO = "r-model-02"

In [13]:
! gcloud artifacts repositories create $PRIVATE_REPO --repository-format=docker \
    --location=$REGION \
    --description="Docker repository for R testing"

Create request issued for: [r-model-02]
Waiting for operation [projects/demos-vertex-ai/locations/us-central1/operation
s/6cd3d644-df78-4b4e-8933-4f37eff94d04] to complete...done.                    
Created repository [r-model-02].


In [14]:
! gcloud artifacts repositories list

Listing items under project demos-vertex-ai, across all locations.

                                                                            ARTIFACT_REGISTRY
REPOSITORY      FORMAT  MODE                 DESCRIPTION                      LOCATION     LABELS  ENCRYPTION          CREATE_TIME          UPDATE_TIME          SIZE (MB)
asia.gcr.io     DOCKER  STANDARD_REPOSITORY                                   asia                 Google-managed key  2022-10-18T19:20:32  2022-10-18T19:20:32  0
automl-beans    DOCKER  STANDARD_REPOSITORY  Docker repository                us-central1          Google-managed key  2022-10-19T20:07:45  2022-10-28T19:53:13  2939.842
eu.gcr.io       DOCKER  STANDARD_REPOSITORY                                   europe               Google-managed key  2022-10-18T19:20:25  2022-10-18T19:20:25  0
gcr.io          DOCKER  STANDARD_REPOSITORY                                   us                   Google-managed key  2022-10-18T19:20:17  2022-10-18T19:20:17  0
my-docke

In [15]:
! gcloud auth configure-docker $REGION-docker.pkg.dev --quiet


{
  "credHelpers": {
    "gcr.io": "gcloud",
    "us.gcr.io": "gcloud",
    "eu.gcr.io": "gcloud",
    "asia.gcr.io": "gcloud",
    "staging-k8s.gcr.io": "gcloud",
    "marketplace.gcr.io": "gcloud"
  }
}
Adding credentials for: us-central1-docker.pkg.dev
Docker configuration file updated.


## Create script 

References: 

* https://github.com/statmike/vertex-ai-mlops/blob/main/08%20-%20R/08b%20Training%20Job%20-%20Vertex%20AI%20Custom%20Model%20-%20R%20-%20Training%20Pipeline%20With%20Custom%20Container.ipynb
* https://github.com/jchavezar/vertex-ai-mlops/blob/main/03%20Tensorflow/03tb%20-%20tfkeras_customjob_xai_tabclass.ipynb

In [27]:
%%writefile src/print.R
#!/usr/bin/env Rscript
# inputs
args <- commandArgs(trailingOnly = TRUE)
project_id <- args[1]
region <- args[2]
gcs_uri_in <- args[3]
gcs_uri_out <- args[4]

cat("project_id: ", project_id)
cat("region: ", region)
cat("gcs_uri_in: ", gcs_uri_in)
cat("gcs_uri_out: ", gcs_uri_out)

# printing test
print("This is printed using the print function")

message("This is printed using the message function")

cat("This is printed using the cat function")

write("This is printed using the write function", stdout())

# view environment vars
Sys.getenv()

Overwriting src/print.R


## Build and push container image

### Create Dockerfile

The docker file for your custom container is built on top of the Deep Learning container -- the same container that is also used for Vertex AI Workbench. In addition, you add two R scripts for model training and serving, respectively.

In [28]:
IMAGE_NAME = "r-model-02"  # @param {type:"string"}
IMAGE_TAG = "latest"  # @param {type:"string"}
IMAGE_URI = f"{REGION}-docker.pkg.dev/{PROJECT_ID}/{PRIVATE_REPO}/{IMAGE_NAME}:{IMAGE_TAG}"
IMAGE_URI

'us-central1-docker.pkg.dev/demos-vertex-ai/r-model-02/r-model-02:latest'

In [29]:
%%writefile ./src/Dockerfile

FROM gcr.io/deeplearning-platform-release/r-cpu.4-1:latest

WORKDIR /root

COPY print.R /root/print.R

RUN apt-get update
RUN apt-get install gfortran -yy

Overwriting ./src/Dockerfile


### Build the Docker container

Next, you build the Docker container image on Cloud Build -- the serverless CI/CD platform.

*Note:* Building the Docker container image may take 5-10 minutes.

In [36]:
! gcloud builds submit --region=$REGION --tag=$IMAGE_URI --timeout=1h ./src --async

Creating temporary tarball archive of 2 file(s) totalling 705 bytes before compression.
Uploading tarball of [./src] to [gs://demos-vertex-ai_cloudbuild/source/1668114926.58813-3081029dd13645619581053d1dd2e40e.tgz]
Created [https://cloudbuild.googleapis.com/v1/projects/demos-vertex-ai/locations/us-central1/builds/d37f3f7a-925c-4751-a76a-e8f80237d434].
Logs are available at [ https://console.cloud.google.com/cloud-build/builds;region=us-central1/d37f3f7a-925c-4751-a76a-e8f80237d434?project=746038361521 ].
ID                                    CREATE_TIME                DURATION  SOURCE                                                                                        IMAGES  STATUS
d37f3f7a-925c-4751-a76a-e8f80237d434  2022-11-10T21:15:27+00:00  -         gs://demos-vertex-ai_cloudbuild/source/1668114926.58813-3081029dd13645619581053d1dd2e40e.tgz  -       QUEUED


## Define Configuration

In [32]:
MODEL_NAME = APP_NAME
MODEL_DISPLAY_NAME = f"{MODEL_NAME}"

PIPELINE_NAME = f"{APP_NAME}-pipeline"
PIPELINE_ROOT = f"{BUCKET_NAME}/pipeline_root/{MODEL_NAME}"
GCS_STAGING = f"{BUCKET_NAME}/pipeline_root/{MODEL_NAME}"

IMAGE_URI # defined above

PIPELINE_JSON_SPEC_PATH = './src/r-test-pipeline-spec.json'

WORKING_DIR = f"{PIPELINE_ROOT}/{TIMESTAMP}"
WORKING_DIR

GCS_URI_IN = f"{BUCKET_NAME}/pipeline_root/data/raw.csv"
GCS_URI_OUT = f"{BUCKET_NAME}/pipeline_root/data/clean.csv"

### Set command arguments for passing into pipeleine 


In [33]:
CMDARGS = [
    "--project_id=" + PROJECT_ID,
    "--region=" + REGION,
    "--gcs_uri_in=" + GCS_URI_IN,
    "--gcs_uri_out=" + GCS_URI_OUT,
]

R code using `commandArgs()` does not used named parameters so parse CMDARGS for the R script:

In [34]:
CMDARGS = [c.split('=')[-1] for c in CMDARGS]
CMDARGS

['demos-vertex-ai',
 'us-central1',
 'gs://demos-vertex-ai-aip-r-model-02-20221110202015/pipeline_root/data/raw.csv',
 'gs://demos-vertex-ai-aip-r-model-02-20221110202015/pipeline_root/data/clean.csv']

## Define Pipeline


* Example - Pipeline - https://github.com/JowGarrido/kfp-r-models/blob/main/pipeline_definition.ipynb 
* SDK Vertex - https://cloud.google.com/python/docs/reference/aiplatform/latest/google.cloud.aiplatform.CustomJob
* API Vertex - https://cloud.google.com/vertex-ai/docs/reference/rest/v1/CustomJobSpec
* Pipeline component - https://google-cloud-pipeline-components.readthedocs.io/en/google-cloud-pipeline-components-1.0.26/google_cloud_pipeline_components.v1.custom_job.html


In [37]:
@dsl.pipeline(
    name = PIPELINE_NAME, 
    pipeline_root = PIPELINE_ROOT)
def r_pipeline():
    print_task = (
        CustomTrainingJobOp(
            project = PROJECT_ID,
            location = REGION,
            display_name = "Run R script",
            base_output_directory = WORKING_DIR,
            worker_pool_specs=[
                {
                    "containerSpec": {
                        "imageUri": IMAGE_URI,
                        "command": ["Rscript", "print.R"],
                        "args": CMDARGS,
                    },
                    "replicaCount": "1",
                    "machineSpec": {
                        "machineType": "n1-standard-16"
                    },
                }
            ],
        )
        .set_display_name("Run R script")
        .set_caching_options(False)
    )

In [38]:
compiler.Compiler().compile(pipeline_func = r_pipeline, 
                            package_path = PIPELINE_JSON_SPEC_PATH)



In [None]:
aiplatform.init(project=PROJECT_ID, location=REGION, staging_bucket=BUCKET_NAME)

In [None]:
job = pipeline_jobs.PipelineJob(
    display_name = 'r_print_test',
    template_path = PIPELINE_JSON_SPEC_PATH,
    pipeline_root = PIPELINE_ROOT,
    enable_caching = False
)

job.submit()