mirror of
https://github.com/GoogleCloudPlatform/vertex-ai-samples.git
synced 2026-09-26 14:42:04 +00:00
67 KiB
67 KiB
In [ ]:
# Copyright 2021 Google LLC
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# https://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.In [ ]:
! pip3 install google-cloud-storageIn [ ]:
import os
if not os.getenv("AUTORUN"):
# Automatically restart kernel after installs
import IPython
app = IPython.Application.instance()
app.kernel.do_shutdown(True)In [ ]:
PROJECT_ID = "[your-project-id]" # @param {type:"string"}In [ ]:
if PROJECT_ID == "" or PROJECT_ID is None or PROJECT_ID == "[your-project-id]":
# Get your GCP project id from gcloud
shell_output = !gcloud config list --format 'value(core.project)' 2>/dev/null
PROJECT_ID = shell_output[0]
print("Project ID:", PROJECT_ID)In [ ]:
! gcloud config set project $PROJECT_IDIn [ ]:
REGION = "us-central1" # @param {type: "string"}In [ ]:
from datetime import datetime
TIMESTAMP = datetime.now().strftime("%Y%m%d%H%M%S")In [ ]:
import os
import sys
# If you are running this notebook in Colab, run this cell and follow the
# instructions to authenticate your Google Cloud account. This provides access
# to your Cloud Storage bucket and lets you submit training jobs and prediction
# requests.
# If on AutoML, then don't execute this code
if not os.path.exists("/opt/deeplearning/metadata/env_version"):
if "google.colab" in sys.modules:
from google.colab import auth as google_auth
google_auth.authenticate_user()
# If you are running this tutorial in a notebook locally, replace the string
# below with the path to your service account key and run this cell to
# authenticate your Google Cloud account.
else:
%env GOOGLE_APPLICATION_CREDENTIALS your_path_to_credentials.json
# Log in to your account on Google Cloud
! gcloud auth loginIn [ ]:
BUCKET_NAME = "[your-bucket-name]" # @param {type:"string"}In [ ]:
if BUCKET_NAME == "" or BUCKET_NAME is None or BUCKET_NAME == "[your-bucket-name]":
BUCKET_NAME = PROJECT_ID + "aip-" + TIMESTAMPIn [ ]:
! gsutil mb -l $REGION gs://$BUCKET_NAMEIn [ ]:
! gsutil ls -al gs://$BUCKET_NAMEIn [ ]:
import json
import time
from googleapiclient import discoveryIn [ ]:
# AutoM location root path for your dataset, model and endpoint resources
PARENT = "projects/" + PROJECT_IDIn [ ]:
cloudml = discovery.build("ml", "v1")In [ ]:
%%writefile cifar/Dockerfile
FROM gcr.io/deeplearning-platform-release/tf2-cpu.2-1
WORKDIR /root
WORKDIR /
# Copies the trainer code to the docker image.
COPY trainer /trainer
# Sets up the entry point to invoke the trainer.
ENTRYPOINT ["python", "-m", "trainer.task"]
In [ ]:
# Add package information
! touch cifar/README.md
setup_cfg = "[egg_info]\n\
tag_build =\n\
tag_date = 0"
! echo "$setup_cfg" > cifar/setup.cfg
setup_py = "import setuptools\n\
# Requires TensorFlow Datasets\n\
setuptools.setup(\n\
install_requires=[\n\
'tensorflow_datasets==1.3.0',\n\
],\n\
packages=setuptools.find_packages())"
! echo "$setup_py" > cifar/setup.py
pkg_info = "Metadata-Version: 1.0\n\
Name: Custom Training CIFAR-10\n\
Version: 0.0.0\n\
Summary: Demonstration training script\n\
Home-page: www.google.com\n\
Author: Google\n\
Author-email: aferlitsch@google.com\n\
License: Public\n\
Description: Demo\n\
Platform: Vertex AI"
! echo "$pkg_info" > cifar/PKG-INFO
# Make the training subfolder
! mkdir cifar/trainer
! touch cifar/trainer/__init__.pyIn [ ]:
%%writefile cifar/trainer/task.py
import tensorflow_datasets as tfds
import tensorflow as tf
from tensorflow.python.client import device_lib
import argparse
import os
import sys
tfds.disable_progress_bar()
parser = argparse.ArgumentParser()
parser.add_argument('--model-dir', dest='model_dir',
default='/tmp/saved_model', type=str, help='Model dir.')
parser.add_argument('--lr', dest='lr',
default=0.01, type=float,
help='Learning rate.')
parser.add_argument('--epochs', dest='epochs',
default=10, type=int,
help='Number of epochs.')
parser.add_argument('--steps', dest='steps',
default=200, type=int,
help='Number of steps per epoch.')
parser.add_argument('--distribute', dest='distribute', type=str, default='single',
help='distributed training strategy')
args = parser.parse_args()
print('Python Version = {}'.format(sys.version))
print('TensorFlow Version = {}'.format(tf.__version__))
print('TF_CONFIG = {}'.format(os.environ.get('TF_CONFIG', 'Not found')))
print('DEVICES', device_lib.list_local_devices())
# Single Machine, single compute device
if args.distribute == 'single':
if tf.test.is_gpu_available():
strategy = tf.distribute.OneDeviceStrategy(device="/gpu:0")
else:
strategy = tf.distribute.OneDeviceStrategy(device="/cpu:0")
# Single Machine, multiple compute device
elif args.distribute == 'mirror':
strategy = tf.distribute.MirroredStrategy()
# Multiple Machine, multiple compute device
elif args.distribute == 'multi':
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()
# Multi-worker configuration
print('num_replicas_in_sync = {}'.format(strategy.num_replicas_in_sync))
# Preparing dataset
BUFFER_SIZE = 10000
BATCH_SIZE = 64
def make_datasets_unbatched():
# Scaling CIFAR10 data from (0, 255] to (0., 1.]
def scale(image, label):
image = tf.cast(image, tf.float32)
image /= 255.0
return image, label
datasets, info = tfds.load(name='cifar10',
with_info=True,
as_supervised=True)
return datasets['train'].map(scale).cache().shuffle(BUFFER_SIZE).repeat()
# Build the Keras model
def build_and_compile_cnn_model():
model = tf.keras.Sequential([
tf.keras.layers.Conv2D(32, 3, activation='relu', input_shape=(32, 32, 3)),
tf.keras.layers.MaxPooling2D(),
tf.keras.layers.Conv2D(32, 3, activation='relu'),
tf.keras.layers.MaxPooling2D(),
tf.keras.layers.Flatten(),
tf.keras.layers.Dense(10, activation='softmax')
])
model.compile(
loss=tf.keras.losses.sparse_categorical_crossentropy,
optimizer=tf.keras.optimizers.SGD(learning_rate=args.lr),
metrics=['accuracy'])
return model
# Train the model
NUM_WORKERS = strategy.num_replicas_in_sync
# Here the batch size scales up by number of workers since
# `tf.data.Dataset.batch` expects the global batch size.
GLOBAL_BATCH_SIZE = BATCH_SIZE * NUM_WORKERS
train_dataset = make_datasets_unbatched().batch(GLOBAL_BATCH_SIZE)
with strategy.scope():
# Creation of dataset, and model building/compiling need to be within
# `strategy.scope()`.
model = build_and_compile_cnn_model()
model.fit(x=train_dataset, epochs=args.epochs, steps_per_epoch=args.steps)
model.save(args.model_dir)
In [ ]:
TRAIN_IMAGE = f"gcr.io/{PROJECT_ID}/cifar_migration:v1"
! docker build cifar -t $TRAIN_IMAGE
! docker push $TRAIN_IMAGEIn [ ]:
JOB_NAME = "custom_container_" + TIMESTAMP
TRAINING_INPUTS = {
"scaleTier": "CUSTOM",
"masterType": "n1-standard-4",
"masterConfig": {"imageUri": TRAIN_IMAGE},
"args": [
"--model-dir=" + "gs://{}/{}".format(BUCKET_NAME, JOB_NAME),
"--epochs=" + str(20),
"--steps=" + str(100),
],
"region": REGION,
}
body = {"jobId": JOB_NAME, "trainingInput": TRAINING_INPUTS}
request = cloudml.projects().jobs().create(parent=PARENT)
request.body = json.loads(json.dumps(TRAINING_INPUTS, indent=2))
print(json.dumps(json.loads(request.to_json()), indent=2))
request = cloudml.projects().jobs().create(parent=PARENT, body=body)In [ ]:
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))In [ ]:
# The full unique ID for the custom training job
custom_training_id = f'{PARENT}/jobs/{response["jobId"]}'
# The short numeric ID for the custom training job
custom_training_short_id = response["jobId"]
print(custom_training_id)In [ ]:
request = cloudml.projects().jobs().get(name=custom_training_id)
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))In [ ]:
while True:
response = cloudml.projects().jobs().get(name=custom_training_id).execute()
if response["state"] != "SUCCEEDED":
print("Training job has not completed:", response["state"])
if response["state"] == "FAILED":
break
else:
break
time.sleep(20)
# model artifact output directory on Google Cloud Storage
model_artifact_dir = response["trainingInput"]["args"][0].split("=")[-1]
print("artifact location " + model_artifact_dir)In [ ]:
import tensorflow as tf
model = tf.keras.models.load_model(model_artifact_dir)In [ ]:
CONCRETE_INPUT = "numpy_inputs"
def _preprocess(bytes_input):
decoded = tf.io.decode_jpeg(bytes_input, channels=3)
decoded = tf.image.convert_image_dtype(decoded, tf.float32)
resized = tf.image.resize(decoded, size=(32, 32))
rescale = tf.cast(resized / 255.0, tf.float32)
return rescale
@tf.function(input_signature=[tf.TensorSpec([None], tf.string)])
def preprocess_fn(bytes_inputs):
decoded_images = tf.map_fn(
_preprocess, bytes_inputs, dtype=tf.float32, back_prop=False
)
return {
CONCRETE_INPUT: decoded_images
} # User needs to make sure the key matches model's input
m_call = tf.function(model.call).get_concrete_function(
[tf.TensorSpec(shape=[None, 32, 32, 3], dtype=tf.float32, name=CONCRETE_INPUT)]
)
@tf.function(
input_signature=[tf.TensorSpec([None], tf.string), tf.TensorSpec([None], tf.string)]
)
def serving_fn(bytes_inputs, key):
images = preprocess_fn(bytes_inputs)
prob = m_call(**images)
return {"prediction": prob, "key": key}
tf.saved_model.save(
model, model_artifact_dir, signatures={"serving_default": serving_fn}
)In [ ]:
loaded = tf.saved_model.load(model_artifact_dir)
tensors_specs = list(loaded.signatures["serving_default"].structured_input_signature)
print("Tensors specs:", tensors_specs)
input_name = [v for k, v in tensors_specs[1].items() if k != "key"][0].name
print("Bytes input tensor name:", input_name)In [ ]:
import base64
import json
import cv2
import numpy as np
import tensorflow as tf
(_, _), (x_test, y_test) = tf.keras.datasets.cifar10.load_data()
test_image_1, test_label_1 = x_test[0], y_test[0]
test_image_2, test_label_2 = x_test[1], y_test[1]
cv2.imwrite("tmp1.jpg", (test_image_1 * 255).astype(np.uint8))
cv2.imwrite("tmp2.jpg", (test_image_2 * 255).astype(np.uint8))
gcs_input_uri = "gs://" + BUCKET_NAME + "/" + "test.json"
with tf.io.gfile.GFile(gcs_input_uri, "w") as f:
for img in ["tmp1.jpg", "tmp2.jpg"]:
bytes = tf.io.read_file(img)
b64str = base64.b64encode(bytes.numpy()).decode("utf-8")
f.write(json.dumps({"key": img, input_name: {"b64": b64str}}) + "\n")
! gsutil cat $gcs_input_uriIn [ ]:
body = {
"jobId": "custom_container_pred_" + TIMESTAMP,
"predictionInput": {
"dataFormat": "JSON",
"inputPaths": gcs_input_uri,
"outputPath": "gs://" + f"{BUCKET_NAME}/batch_output/",
"runtime_version": "2.1",
"uri": model_artifact_dir,
"region": REGION,
},
}
request = cloudml.projects().jobs().create(parent=PARENT)
request.body = json.loads(json.dumps(body, indent=2))
print(json.dumps(json.loads(request.to_json()), indent=2))
request = cloudml.projects().jobs().create(parent=PARENT, body=body)In [ ]:
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))In [ ]:
# The full unique ID for the batch prediction job
batch_job_id = PARENT + "/jobs/" + response["jobId"]
print(batch_job_id)In [ ]:
request = cloudml.projects().jobs().get(name=batch_job_id)
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))In [ ]:
while True:
response = request = cloudml.projects().jobs().get(name=batch_job_id).execute()
if response["state"] != "SUCCEEDED":
print("The job has not completed:", response["state"])
if response["state"] == "FAILED":
break
else:
folder = response["predictionInput"]["outputPath"][:-1]
! gsutil ls $folder/prediction*
! gsutil cat $folder/prediction*
break
time.sleep(60)In [ ]:
request = cloudml.projects().models().create(parent=PARENT)
request.body = json.loads(
json.dumps({"name": "custom_container_" + TIMESTAMP}, indent=2)
)
print(json.dumps(json.loads(request.to_json()), indent=2))
request = (
cloudml.projects()
.models()
.create(parent=PARENT, body={"name": "custom_container_" + TIMESTAMP})
)In [ ]:
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))In [ ]:
# The full unique ID for the training pipeline
model_id = response["name"]
# The short numeric ID for the training pipeline
model_short_name = model_id.split("/")[-1]
print(model_id)In [ ]:
version = {
"name": "custom_container_" + TIMESTAMP,
"deploymentUri": model_artifact_dir,
"runtimeVersion": "2.1",
"framework": "TENSORFLOW",
"pythonVersion": "3.7",
"machineType": "mls1-c1-m2",
}
request = cloudml.projects().models().versions().create(parent=response["name"])
request.body = json.loads(json.dumps(version, indent=2))
print(json.dumps(json.loads(request.to_json()), indent=2))
request = (
cloudml.projects().models().versions().create(parent=response["name"], body=version)
)In [ ]:
response = request.execute()In [ ]:
print(json.dumps(response, indent=2))Warning:
Output truncated. This notebook contains too many cells to display efficiently.