### SETUP

In [151]:
!pip install -U tfx
#!pip install --user tfx tensorflow Pillow tensorflow_datasets matplotlib azure-storage-blob object-detection

Defaulting to user installation because normal site-packages is not writeable
You should consider upgrading via the '/usr/local/bin/python -m pip install --upgrade pip' command.[0m[33m
[0m

In [152]:
# Import tensorflow and TFX modules
import tensorflow as tf
print('TensorFlow version: {}'.format(tf.__version__))
from tfx import v1 as tfx
print('TFX version: {}'.format(tfx.__version__))


TensorFlow version: 2.11.0
TFX version: 1.12.0


In [153]:
import os

PIPELINE_NAME = "penguin-pipeline2"

# Output directory to store artifacts generated from the pipeline.
PIPELINE_ROOT = os.path.join('pipelines', PIPELINE_NAME)
# Path to a SQLite DB file to use as an MLMD storage.
METADATA_PATH = os.path.join('metadata', PIPELINE_NAME, 'metadata.db')
# Output directory where created models from the pipeline will be exported.
SERVING_MODEL_DIR = os.path.join('serving_model', PIPELINE_NAME)

from absl import logging
logging.set_verbosity(logging.INFO)  # Set default logging level.


### PREP DATA

In [154]:
import urllib.request
import tempfile

#processed data file

DATA_ROOT = os.path.join('penguin-data')
PROCESSED_DATA_ROOT = os.path.join(DATA_ROOT, 'processed')

os.makedirs(DATA_ROOT + '/processed', exist_ok=True)

_data_path = 'https://raw.githubusercontent.com/tensorflow/tfx/master/tfx/examples/penguin/data/labelled/penguins_processed.csv'
_data_filepath = os.path.join(DATA_ROOT,'processed', 'data-processed.csv')
urllib.request.urlretrieve(_data_path, _data_filepath)

#Non processed data

#DATA_ROOT = os.path.join('penguin-data')
#FULL_DATA_ROOT = os.path.join(DATA_ROOT, 'full-set')
#INCOMPLETE_DATA_ROOT = os.path.join(DATA_ROOT, 'incomplete-set')

#os.makedirs(DATA_ROOT + '/full-set', exist_ok=True)
#_data_path = 'https://storage.googleapis.com/download.tensorflow.org/data/palmer_penguins/penguins_size.csv'
#_data_filepath = os.path.join(DATA_ROOT,'full-set', 'data-full.csv')
#urllib.request.urlretrieve(_data_path, _data_filepath)


('penguin-data/processed/data-processed.csv',
 <http.client.HTTPMessage at 0x7fece76e4390>)

In [155]:
!head {_data_filepath}


species,culmen_length_mm,culmen_depth_mm,flipper_length_mm,body_mass_g
0,0.2545454545454545,0.6666666666666666,0.15254237288135594,0.2916666666666667
0,0.26909090909090905,0.5119047619047618,0.23728813559322035,0.3055555555555556
0,0.29818181818181805,0.5833333333333334,0.3898305084745763,0.1527777777777778
0,0.16727272727272732,0.7380952380952381,0.3559322033898305,0.20833333333333334
0,0.26181818181818167,0.892857142857143,0.3050847457627119,0.2638888888888889
0,0.24727272727272717,0.5595238095238096,0.15254237288135594,0.2569444444444444
0,0.25818181818181823,0.773809523809524,0.3898305084745763,0.5486111111111112
0,0.32727272727272727,0.5357142857142859,0.1694915254237288,0.1388888888888889
0,0.23636363636363636,0.9642857142857142,0.3220338983050847,0.3055555555555556


### DELETE ENTRIES WITH NA FIELDS

In [156]:
#!sed -i '/\bNA\b/d' {_data_filepath}
#!head {_data_filepath}


### Create multiple datasets

In [157]:
## Remove 50 entries of Chinstrap species
#%run remove_csv.py

### CREATE DATA.CSV file with incomplete dataset

### CREATING SCHEMA

In [158]:
import shutil

SCHEMA_PATH = 'schema'

_schema_uri = 'https://raw.githubusercontent.com/tensorflow/tfx/master/tfx/examples/penguin/schema/raw/schema.pbtxt'
_schema_filename = 'schema.pbtxt'
_schema_filepath = os.path.join(SCHEMA_PATH, _schema_filename)

os.makedirs(SCHEMA_PATH, exist_ok=True)
urllib.request.urlretrieve(_schema_uri, _schema_filepath)


('schema/schema.pbtxt', <http.client.HTTPMessage at 0x7fece76dad90>)

### CREATE FILES FOR COMPONENT FUNCTIONS

In [159]:
_module_file = 'penguin_utils.py'


In [160]:
%%writefile {_module_file}


from typing import List, Text
from absl import logging
import tensorflow as tf
from tensorflow import keras
from tensorflow_metadata.proto.v0 import schema_pb2
import tensorflow_transform as tft
from tensorflow_transform.tf_metadata import schema_utils

from tfx import v1 as tfx
from tfx_bsl.public import tfxio

# Specify features that we will use.
_FEATURE_KEYS = [
    'culmen_length_mm', 'culmen_depth_mm', 'flipper_length_mm', 'body_mass_g'
]
_LABEL_KEY = 'species'

_TRAIN_BATCH_SIZE = 20
_EVAL_BATCH_SIZE = 10


# NEW: TFX Transform will call this function.
def preprocessing_fn(inputs):
  """tf.transform's callback function for preprocessing inputs.

  Args:
    inputs: map from feature keys to raw not-yet-transformed features.

  Returns:
    Map from string feature key to transformed feature.
  """
  outputs = {}

  # Uses features defined in _FEATURE_KEYS only.
  for key in _FEATURE_KEYS:
    # tft.scale_to_z_score computes the mean and variance of the given feature
    # and scales the output based on the result.
    outputs[key] = tft.scale_to_z_score(inputs[key])

  # For the label column we provide the mapping from string to index.
  # We could instead use `tft.compute_and_apply_vocabulary()` in order to
  # compute the vocabulary dynamically and perform a lookup.
  # Since in this example there are only 3 possible values, we use a hard-coded
  # table for simplicity.
  table_keys = ['Adelie', 'Chinstrap', 'Gentoo']
  initializer = tf.lookup.KeyValueTensorInitializer(
      keys=table_keys,
      values=tf.cast(tf.range(len(table_keys)), tf.int64),
      key_dtype=tf.string,
      value_dtype=tf.int64)
  table = tf.lookup.StaticHashTable(initializer, default_value=-1)
  outputs[_LABEL_KEY] = table.lookup(inputs[_LABEL_KEY])

  return outputs


# NEW: This function will apply the same transform operation to training data
#      and serving requests.
def _apply_preprocessing(raw_features, tft_layer):
  transformed_features = tft_layer(raw_features)
  if _LABEL_KEY in raw_features:
    transformed_label = transformed_features.pop(_LABEL_KEY)
    return transformed_features, transformed_label
  else:
    return transformed_features, None


# NEW: This function will create a handler function which gets a serialized
#      tf.example, preprocess and run an inference with it.
def _get_serve_tf_examples_fn(model, tf_transform_output):
  # We must save the tft_layer to the model to ensure its assets are kept and
  # tracked.
  model.tft_layer = tf_transform_output.transform_features_layer()

  @tf.function(input_signature=[
      tf.TensorSpec(shape=[None], dtype=tf.string, name='examples')
  ])
  def serve_tf_examples_fn(serialized_tf_examples):
    # Expected input is a string which is serialized tf.Example format.
    feature_spec = tf_transform_output.raw_feature_spec()
    # Because input schema includes unnecessary fields like 'species' and
    # 'island', we filter feature_spec to include required keys only.
    required_feature_spec = {
        k: v for k, v in feature_spec.items() if k in _FEATURE_KEYS
    }
    parsed_features = tf.io.parse_example(serialized_tf_examples,
                                          required_feature_spec)

    # Preprocess parsed input with transform operation defined in
    # preprocessing_fn().
    transformed_features, _ = _apply_preprocessing(parsed_features,
                                                   model.tft_layer)
    # Run inference with ML model.
    return model(transformed_features)

  return serve_tf_examples_fn


def _input_fn(file_pattern: List[Text],
              data_accessor: tfx.components.DataAccessor,
              tf_transform_output: tft.TFTransformOutput,
              batch_size: int = 200) -> tf.data.Dataset:
  """Generates features and label for tuning/training.

  Args:
    file_pattern: List of paths or patterns of input tfrecord files.
    data_accessor: DataAccessor for converting input to RecordBatch.
    tf_transform_output: A TFTransformOutput.
    batch_size: representing the number of consecutive elements of returned
      dataset to combine in a single batch

  Returns:
    A dataset that contains (features, indices) tuple where features is a
      dictionary of Tensors, and indices is a single Tensor of label indices.
  """
  dataset = data_accessor.tf_dataset_factory(
      file_pattern,
      tfxio.TensorFlowDatasetOptions(batch_size=batch_size),
      schema=tf_transform_output.raw_metadata.schema)

  transform_layer = tf_transform_output.transform_features_layer()
  def apply_transform(raw_features):
    return _apply_preprocessing(raw_features, transform_layer)

  return dataset.map(apply_transform).repeat()


def _build_keras_model() -> tf.keras.Model:
  """Creates a DNN Keras model for classifying penguin data.

  Returns:
    A Keras Model.
  """
  # The model below is built with Functional API, please refer to
  # https://www.tensorflow.org/guide/keras/overview for all API options.
  inputs = [
      keras.layers.Input(shape=(1,), name=key)
      for key in _FEATURE_KEYS
  ]
  d = keras.layers.concatenate(inputs)
  for _ in range(2):
    d = keras.layers.Dense(8, activation='relu')(d)
  outputs = keras.layers.Dense(3)(d)

  model = keras.Model(inputs=inputs, outputs=outputs)
  model.compile(
      optimizer=keras.optimizers.Adam(1e-2),
      loss=tf.keras.losses.SparseCategoricalCrossentropy(from_logits=True),
      metrics=[keras.metrics.SparseCategoricalAccuracy()])

  model.summary(print_fn=logging.info)
  return model


# TFX Trainer will call this function.
def run_fn(fn_args: tfx.components.FnArgs):
  """Train the model based on given args.

  Args:
    fn_args: Holds args used to train the model as name/value pairs.
  """
  tf_transform_output = tft.TFTransformOutput(fn_args.transform_output)

  train_dataset = _input_fn(
      fn_args.train_files,
      fn_args.data_accessor,
      tf_transform_output,
      batch_size=_TRAIN_BATCH_SIZE)
  eval_dataset = _input_fn(
      fn_args.eval_files,
      fn_args.data_accessor,
      tf_transform_output,
      batch_size=_EVAL_BATCH_SIZE)

  model = _build_keras_model()
  model.fit(
      train_dataset,
      steps_per_epoch=fn_args.train_steps,
      validation_data=eval_dataset,
      validation_steps=fn_args.eval_steps)

  # NEW: Save a computation graph including transform layer.
  signatures = {
      'serving_default': _get_serve_tf_examples_fn(model, tf_transform_output),
  }
  model.save(fn_args.serving_model_dir, save_format='tf', signatures=signatures)


Overwriting penguin_utils.py


In [161]:
%%writefile {_module_file}

# Copied from https://www.tensorflow.org/tfx/tutorials/tfx/penguin_simple

from typing import List
from absl import logging
import tensorflow as tf
from tensorflow import keras
from tensorflow_transform.tf_metadata import schema_utils

from tfx.components.trainer.executor import TrainerFnArgs
from tfx.components.trainer.fn_args_utils import DataAccessor
from tfx_bsl.tfxio import dataset_options
from tensorflow_metadata.proto.v0 import schema_pb2

_FEATURE_KEYS = [
    'culmen_length_mm', 'culmen_depth_mm', 'flipper_length_mm', 'body_mass_g'
]
_LABEL_KEY = 'species'

_TRAIN_BATCH_SIZE = 20
_EVAL_BATCH_SIZE = 10

# Since we're not generating or creating a schema, we will instead create
# a feature spec.  Since there are a fairly small number of features this is
# manageable for this dataset.
_FEATURE_SPEC = {
    **{
        feature: tf.io.FixedLenFeature(shape=[1], dtype=tf.float32)
           for feature in _FEATURE_KEYS
       },
    _LABEL_KEY: tf.io.FixedLenFeature(shape=[1], dtype=tf.int64)
}


def _input_fn(file_pattern: List[str],
              data_accessor: DataAccessor,
              schema: schema_pb2.Schema,
              batch_size: int = 200) -> tf.data.Dataset:
  """Generates features and label for training.

  Args:
    file_pattern: List of paths or patterns of input tfrecord files.
    data_accessor: DataAccessor for converting input to RecordBatch.
    schema: schema of the input data.
    batch_size: representing the number of consecutive elements of returned
      dataset to combine in a single batch

  Returns:
    A dataset that contains (features, indices) tuple where features is a
      dictionary of Tensors, and indices is a single Tensor of label indices.
  """
  return data_accessor.tf_dataset_factory(
      file_pattern,
      dataset_options.TensorFlowDatasetOptions(
          batch_size=batch_size, label_key=_LABEL_KEY),
      schema=schema).repeat()


def _build_keras_model() -> tf.keras.Model:
  """Creates a DNN Keras model for classifying penguin data.

  Returns:
    A Keras Model.
  """
  # The model below is built with Functional API, please refer to
  # https://www.tensorflow.org/guide/keras/overview for all API options.
  inputs = [keras.layers.Input(shape=(1,), name=f) for f in _FEATURE_KEYS]
  d = keras.layers.concatenate(inputs)
  for _ in range(2):
    d = keras.layers.Dense(8, activation='relu')(d)
  outputs = keras.layers.Dense(3)(d)

  model = keras.Model(inputs=inputs, outputs=outputs)
  model.compile(
      optimizer=keras.optimizers.Adam(1e-2),
      loss=tf.keras.losses.SparseCategoricalCrossentropy(from_logits=True),
      metrics=[keras.metrics.SparseCategoricalAccuracy()])

  model.summary(print_fn=logging.info)
  return model


# TFX Trainer will call this function.
def run_fn(fn_args: TrainerFnArgs):
  """Train the model based on given args.

  Args:
    fn_args: Holds args used to train the model as name/value pairs.
  """

  # This schema is usually either an output of SchemaGen or a manually-curated
  # version provided by pipeline author. A schema can also derived from TFT
  # graph if a Transform component is used. In the case when either is missing,
  # `schema_from_feature_spec` could be used to generate schema from very simple
  # feature_spec, but the schema returned would be very primitive.
  schema = schema_utils.schema_from_feature_spec(_FEATURE_SPEC)
  
  train_dataset = _input_fn(
      fn_args.train_files,
      fn_args.data_accessor,
      schema,
      batch_size=_TRAIN_BATCH_SIZE)
  eval_dataset = _input_fn(
      fn_args.eval_files,
      fn_args.data_accessor,
      schema,
      batch_size=_EVAL_BATCH_SIZE)

  model = _build_keras_model()
  model.fit(
      train_dataset,
      steps_per_epoch=fn_args.train_steps,
      validation_data=eval_dataset,
      validation_steps=fn_args.eval_steps)

  # The result of the training should be saved in `fn_args.serving_model_dir`
  # directory.
  model.save(fn_args.serving_model_dir, save_format='tf')


Overwriting penguin_utils.py


### PIPELINE DEFINITION

In [188]:
import ml_metadata as mlmd
from ml_metadata.metadata_store import metadata_store
from ml_metadata.proto import metadata_store_pb2
from tfx.types.channel_utils import external_project_artifact_query
import pprint

"""
def create_custom_model_artifact(model_path):
    model_artifact = tfx.types.standard_artifacts.Model()
    model_artifact.uri = model_path
    model_artifact.pipeline_name = 'penguin-pipeline1'
    model_artifact.set_int_custom_property('artifact_id', 999999) # Use a unique artifact ID
    return model_artifact
"""
# Define a custom filter function to filter out the custom model artifact
#def filter_baseline_model(artifact: metadata_store_pb2.Artifact):
#    return artifact.mlmd_artifact.custom_properties['artifact_id'].int_value != 999999

#model_path = 'pipelines/penguin-pipeline1/Trainer/model/5'
#print(create_custom_model_artifact(model_path))
#print(filter_baseline_model(create_custom_model_artifact('pipelines/penguin-pipeline1/Trainer/model/2')))
 
# Create a custom Model artifact pointing to the saved_model.pb file
#custom_model_artifact = create_custom_model_artifact(model_path)

# Create a Model channel pointing to the custom Model artifact
baseline_model_channel = tfx.dsl.Channel(type=tfx.types.standard_artifacts.Model)
baseline_model_channel.artifacts = ["Hej svej go", "brotha"]
baseline_model_channel.additional_properties = {'hej': 'svej'}
print(baseline_model_channel)

external_pipeline_channel = external_project_artifact_query(
    artifact_type=tfx.types.standard_artifacts.Model,
    pipeline_name='penguin-pipeline1',
    pipeline_run_id='2023-04-27T09:11:05.045733',
    producer_component_id='Trainer', 
    output_key='model',
    project_name="penguin_pipeline1",
    project_owner="penguin_pipeline1",
    mlmd_service_target="metadata/penguin-pipeline1"
    )

print('EXTERNAL CHANNEL: ', external_pipeline_channel)



Channel(
    type_name: Model
    artifacts: []
    additional_properties: {'hej': 'svej'}
    additional_custom_properties: {}
)
EXTERNAL CHANNEL:  ExternalProjectChannel(project_owner=penguin_pipeline1, project_name=penguin_pipeline1, mlmd_service_target=metadata/penguin-pipeline1, pipeline_name=penguin-pipeline1, producer_component_id=Trainer, output_key=model, pipeline_run_id=2023-04-27T09:11:05.045733)


In [186]:
import tensorflow_model_analysis as tfma

def _create_pipeline(pipeline_name: str, pipeline_root: str, data_root: str,
                      module_file: str, serving_model_dir: str,
                     metadata_path: str) -> tfx.dsl.Pipeline:
    """Implements the penguin pipeline with TFX."""
    # Brings data into the pipeline or otherwise joins/converts training data.
    example_gen = tfx.components.CsvExampleGen(input_base=data_root)

    # Uses user-provided Python function that trains a model.
    trainer = tfx.components.Trainer(
        module_file=module_file,
        examples=example_gen.outputs['examples'],

        # NEW: Pass transform_graph to the trainer.
        #transform_graph=transform.outputs['transform_graph'],

        train_args=tfx.proto.TrainArgs(num_steps=100),
        eval_args=tfx.proto.EvalArgs(num_steps=5))


# we want evaluator functionality here
  # NEW: Uses TFMA to compute evaluation statistics over features of a model and
  #   perform quality validation of a candidate model (compared to a baseline).

    eval_config = tfma.EvalConfig(
      model_specs=[tfma.ModelSpec(label_key='species')],
      slicing_specs=[
          # An empty slice spec means the overall slice, i.e. the whole dataset.
          tfma.SlicingSpec(),
          # Calculate metrics for each penguin species.
          tfma.SlicingSpec(feature_keys=['species']),
          ],
      metrics_specs=[
          tfma.MetricsSpec(per_slice_thresholds={
              'sparse_categorical_accuracy':
                  tfma.PerSliceMetricThresholds(thresholds=[
                      tfma.PerSliceMetricThreshold(
                          slicing_specs=[tfma.SlicingSpec()],
                          threshold=tfma.MetricThreshold(
                              value_threshold=tfma.GenericValueThreshold(
                                   lower_bound={'value': 0.6}),
                              # Change threshold will be ignored if there is no
                              # baseline model resolved from MLMD (first run).
                              change_threshold=tfma.GenericChangeThreshold(
                                  direction=tfma.MetricDirection.HIGHER_IS_BETTER,
                                  absolute={'value': -1e-10}))
                       )]),
          })],
      )


        
    # NEW: Get the latest blessed model for Evaluator.
    model_resolver = tfx.dsl.Resolver(
        strategy_class=tfx.dsl.experimental.LatestBlessedModelStrategy,
        model=tfx.dsl.Channel(type=tfx.types.standard_artifacts.Model),
        model_blessing=tfx.dsl.Channel(
            type=tfx.types.standard_artifacts.ModelBlessing)).with_id(
        'latest_blessed_model_resolver')
    
    evaluator = tfx.components.Evaluator(
        examples=example_gen.outputs['examples'],
        model=trainer.outputs['model'],
        baseline_model=external_pipeline_channel,
        #baseline_model=model_resolver.outputs['model'],
        # Change threshold will be ignored if there is no baseline (first run).
        eval_config=eval_config)

    # Pushes the model to a filesystem destination.
    pusher = tfx.components.Pusher(
        model=trainer.outputs['model'],
        model_blessing=evaluator.outputs['blessing'],
        push_destination=tfx.proto.PushDestination(
            filesystem=tfx.proto.PushDestination.Filesystem(
                base_directory=serving_model_dir)))

    components = [
        example_gen,
        #statistics_gen,
        #schema_importer,
        #example_validator,

        #transform,  # NEW: Transform component was added to the pipeline.

        trainer,
        model_resolver,
        evaluator,  #evaluator
        pusher,
    ]

    return tfx.dsl.Pipeline(
        pipeline_name=pipeline_name,
        pipeline_root=pipeline_root,
        metadata_connection_config=tfx.orchestration.metadata
        .sqlite_metadata_connection_config(metadata_path),
        components=components)


In [187]:
#change to full data instead of incomplete
tfx.orchestration.LocalDagRunner().run(
_create_pipeline(
    pipeline_name=PIPELINE_NAME,
    pipeline_root=PIPELINE_ROOT,
    data_root=PROCESSED_DATA_ROOT,
    #schema_path=SCHEMA_PATH,
    module_file=_module_file,
    serving_model_dir=SERVING_MODEL_DIR,
    metadata_path=METADATA_PATH))

INFO:absl:Generating ephemeral wheel package for '/workspaces/tfx_test_case/penguin_utils.py' (including modules: ['penguin_utils', 'remove_csv']).
INFO:absl:User module package has hash fingerprint version 2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0.
INFO:absl:Executing: ['/usr/local/bin/python', '/tmp/tmptmahof1o/_tfx_generated_setup.py', 'bdist_wheel', '--bdist-dir', '/tmp/tmplduxcr7r', '--dist-dir', '/tmp/tmpohm6ospn']


running bdist_wheel
running build
running build_py
creating build
creating build/lib
copying penguin_utils.py -> build/lib
copying remove_csv.py -> build/lib
installing to /tmp/tmplduxcr7r
running install
running install_lib
copying build/lib/penguin_utils.py -> /tmp/tmplduxcr7r
copying build/lib/remove_csv.py -> /tmp/tmplduxcr7r
running install_egg_info


INFO:absl:Successfully built user code wheel distribution at 'pipelines/penguin-pipeline2/_wheels/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3-none-any.whl'; target user module is 'penguin_utils'.
INFO:absl:Full user module path is 'penguin_utils@pipelines/penguin-pipeline2/_wheels/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3-none-any.whl'
INFO:absl:Using deployment config:
 executor_specs {
  key: "CsvExampleGen"
  value {
    beam_executable_spec {
      python_executor_spec {
        class_path: "tfx.components.example_gen.csv_example_gen.executor.Executor"
      }
    }
  }
}
executor_specs {
  key: "Evaluator"
  value {
    beam_executable_spec {
      python_executor_spec {
        class_path: "tfx.components.evaluator.executor.Executor"
      }
    }
  }
}
executor_specs {
  key: "Pusher"
  value {
    python_class_executable_spec {
      class_path: "tfx.components.pusher.executo

running egg_info
creating tfx_user_code_Trainer.egg-info
writing tfx_user_code_Trainer.egg-info/PKG-INFO
writing dependency_links to tfx_user_code_Trainer.egg-info/dependency_links.txt
writing top-level names to tfx_user_code_Trainer.egg-info/top_level.txt
writing manifest file 'tfx_user_code_Trainer.egg-info/SOURCES.txt'
reading manifest file 'tfx_user_code_Trainer.egg-info/SOURCES.txt'
writing manifest file 'tfx_user_code_Trainer.egg-info/SOURCES.txt'
Copying tfx_user_code_Trainer.egg-info to /tmp/tmplduxcr7r/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3.7.egg-info
running install_scripts
creating /tmp/tmplduxcr7r/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0.dist-info/WHEEL
creating '/tmp/tmpohm6ospn/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3-none-any.whl' and adding '/tmp/tmplduxcr7r' to it
adding 'penguin_utils.py'
adding 'remove_csv.py'

INFO:absl:select span and version = (0, None)
INFO:absl:latest span and version = (0, None)
INFO:absl:MetadataStore with DB connection initialized
INFO:absl:Going to run a new execution 43
INFO:absl:Going to run a new execution: ExecutionInfo(execution_id=43, input_dict={}, output_dict=defaultdict(<class 'list'>, {'examples': [Artifact(artifact: uri: "pipelines/penguin-pipeline2/CsvExampleGen/examples/43"
custom_properties {
  key: "input_fingerprint"
  value {
    string_value: "split:single_split,num_files:1,total_bytes:25648,xor_checksum:1682589380,sum_checksum:1682589380"
  }
}
custom_properties {
  key: "span"
  value {
    int_value: 0
  }
}
, artifact_type: name: "Examples"
properties {
  key: "span"
  value: INT
}
properties {
  key: "split_names"
  value: STRING
}
properties {
  key: "version"
  value: INT
}
base_type: DATASET
)]}), exec_properties={'input_config': '{\n  "splits": [\n    {\n      "name": "single_split",\n      "pattern": "*"\n    }\n  ]\n}', 'input_base': 'pen

Processing ./pipelines/penguin-pipeline2/_wheels/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3-none-any.whl


You should consider upgrading via the '/usr/local/bin/python -m pip install --upgrade pip' command.
INFO:absl:Successfully installed 'pipelines/penguin-pipeline2/_wheels/tfx_user_code_Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0-py3-none-any.whl'.
INFO:absl:Training model.
INFO:absl:Feature body_mass_g has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_depth_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature flipper_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature species has a shape dim {
  size: 1
}
. Setting to DenseTensor.


Installing collected packages: tfx-user-code-Trainer
Successfully installed tfx-user-code-Trainer-0.0+2e96f662907f75f30f530f04b3aab5018be8a1ca9ae45f8f8bc0c97ab53843d0


INFO:absl:Feature body_mass_g has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_depth_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature flipper_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature species has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature body_mass_g has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_depth_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature flipper_length_mm has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature species has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature body_mass_g has a shape dim {
  size: 1
}
. Setting to DenseTensor.
INFO:absl:Feature culmen_depth_m





INFO:tensorflow:Assets written to: pipelines/penguin-pipeline2/Trainer/model/45/Format-Serving/assets


INFO:tensorflow:Assets written to: pipelines/penguin-pipeline2/Trainer/model/45/Format-Serving/assets
INFO:absl:Training complete. Model written to pipelines/penguin-pipeline2/Trainer/model/45/Format-Serving. ModelRun written to pipelines/penguin-pipeline2/Trainer/model_run/45
INFO:absl:Cleaning up stateless execution info.
INFO:absl:Execution 45 succeeded.
INFO:absl:Cleaning up stateful execution info.
INFO:absl:Publishing output artifacts defaultdict(<class 'list'>, {'model_run': [Artifact(artifact: uri: "pipelines/penguin-pipeline2/Trainer/model_run/45"
, artifact_type: name: "ModelRun"
)], 'model': [Artifact(artifact: uri: "pipelines/penguin-pipeline2/Trainer/model/45"
, artifact_type: name: "Model"
base_type: MODEL
)]}) for execution 45
INFO:absl:MetadataStore with DB connection initialized
INFO:absl:Component Trainer is finished.
INFO:absl:Component Evaluator is running.
INFO:absl:Running launcher for node_info {
  type {
    name: "tfx.components.evaluator.component.Evaluator"
 

In [181]:
from ml_metadata.proto import metadata_store_pb2
# Non-public APIs, just for showcase.
from tfx.orchestration.portable.mlmd import execution_lib

# TODO(b/171447278): Move these functions into the TFX library.

def get_latest_artifacts(metadata, pipeline_name, component_id):
  """Output artifacts of the latest run of the component."""
  context = metadata.store.get_context_by_type_and_name(
      'node', f'{pipeline_name}.{component_id}')
  #print(context)
  executions = metadata.store.get_executions_by_context(context.id)
  latest_execution = max(executions,
                         key=lambda e:e.last_update_time_since_epoch)
  return execution_lib.get_output_artifacts(metadata, latest_execution.id)


In [182]:
# Non-public APIs, just for showcase.
from tfx.orchestration.metadata import Metadata
from tfx.types import standard_component_specs

metadata_connection_config = tfx.orchestration.metadata.sqlite_metadata_connection_config(
    METADATA_PATH)

with Metadata(metadata_connection_config) as metadata_handler:
  # Find output artifacts from MLMD.
  evaluator_output = get_latest_artifacts(metadata_handler, PIPELINE_NAME,
                                          'Evaluator')
  eval_artifact = evaluator_output[standard_component_specs.EVALUATION_KEY][0]
  #print(eval_artifact)

INFO:absl:MetadataStore with DB connection initialized


In [183]:
import tensorflow_model_analysis as tfma
#print(eval_artifact.uri)
eval_result = tfma.load_eval_result(eval_artifact.uri)
print(eval_result)
tfma.view.render_slicing_metrics(eval_result, slicing_column='species')


EvalResult(slicing_metrics=[((('species', 0),), {'': {'': {'sparse_categorical_accuracy': {'doubleValue': 0.978723406791687}, 'loss': {'doubleValue': 0.12361576408147812}}}}), ((), {'': {'': {'sparse_categorical_accuracy': {'doubleValue': 0.9599999785423279}, 'loss': {'doubleValue': 0.13475194573402405}}}}), ((('species', 1),), {'': {'': {'sparse_categorical_accuracy': {'doubleValue': 0.8636363744735718}, 'loss': {'doubleValue': 0.33218708634376526}}}}), ((('species', 2),), {'': {'': {'sparse_categorical_accuracy': {'doubleValue': 1.0}, 'loss': {'doubleValue': 0.011520599946379662}}}})], plots=[((('species', 0),), None), ((), None), ((('species', 1),), None), ((('species', 2),), None)], attributions=[((('species', 0),), None), ((), None), ((('species', 1),), None), ((('species', 2),), None)], config=model_specs {
  label_key: "species"
}
slicing_specs {
}
slicing_specs {
  feature_keys: "species"
}
metrics_specs {
  model_names: ""
  per_slice_thresholds {
    key: "sparse_categorical_

SlicingMetricsViewer(config={'weightedExamplesColumn': 'example_count'}, data=[{'slice': 'species:0', 'metrics…

In [168]:
# List files in created model directory.
!find {SERVING_MODEL_DIR}


serving_model/penguin-pipeline2
serving_model/penguin-pipeline2/1682588571
serving_model/penguin-pipeline2/1682588571/keras_metadata.pb
serving_model/penguin-pipeline2/1682588571/saved_model.pb
serving_model/penguin-pipeline2/1682588571/assets
serving_model/penguin-pipeline2/1682588571/variables
serving_model/penguin-pipeline2/1682588571/variables/variables.index
serving_model/penguin-pipeline2/1682588571/variables/variables.data-00000-of-00001
serving_model/penguin-pipeline2/1682588571/fingerprint.pb
serving_model/penguin-pipeline2/1682589392
serving_model/penguin-pipeline2/1682589392/keras_metadata.pb
serving_model/penguin-pipeline2/1682589392/saved_model.pb
serving_model/penguin-pipeline2/1682589392/assets
serving_model/penguin-pipeline2/1682589392/variables
serving_model/penguin-pipeline2/1682589392/variables/variables.index
serving_model/penguin-pipeline2/1682589392/variables/variables.data-00000-of-00001
serving_model/penguin-pipeline2/1682589392/fingerprint.pb
serving_model/peng

In [169]:
!saved_model_cli show --dir {SERVING_MODEL_DIR}/$(ls -1 {SERVING_MODEL_DIR} | sort -nr | head -1) --tag_set serve --signature_def serving_default


2023-04-27 09:56:33.915685: I tensorflow/core/platform/cpu_feature_guard.cc:193] This TensorFlow binary is optimized with oneAPI Deep Neural Network Library (oneDNN) to use the following CPU instructions in performance-critical operations:  AVX2 FMA
To enable them in other operations, rebuild TensorFlow with the appropriate compiler flags.
2023-04-27 09:56:34.205684: W tensorflow/compiler/xla/stream_executor/platform/default/dso_loader.cc:64] Could not load dynamic library 'libcudart.so.11.0'; dlerror: libcudart.so.11.0: cannot open shared object file: No such file or directory
2023-04-27 09:56:34.205734: I tensorflow/compiler/xla/stream_executor/cuda/cudart_stub.cc:29] Ignore above cudart dlerror if you do not have a GPU set up on your machine.
2023-04-27 09:56:35.406510: W tensorflow/compiler/xla/stream_executor/platform/default/dso_loader.cc:64] Could not load dynamic library 'libnvinfer.so.7'; dlerror: libnvinfer.so.7: cannot open shared object file: No such file or directory
2023-

In [170]:
# Find a model with the latest and oldest timestamp.
model_dirs = (item for item in os.scandir(SERVING_MODEL_DIR) if item.is_dir())
print('model_dirs ', (item for item in os.scandir(SERVING_MODEL_DIR) if item.is_dir()))
model_path_new = max(model_dirs, key=lambda i: int(i.name)).path

model_dirs = (item for item in os.scandir(SERVING_MODEL_DIR) if item.is_dir())
model_path_old = min(model_dirs, key=lambda i: int(i.name)).path
print('max value ', model_path_new, ' min value ', model_path_old)
loaded_model_new = tf.keras.models.load_model(model_path_new)
loaded_model_old = tf.keras.models.load_model(model_path_old)
inference_fn_new = loaded_model_new.signatures['serving_default']
inference_fn_old = loaded_model_old.signatures['serving_default']


model_dirs  <generator object <genexpr> at 0x7fece76ea5d0>
max value  serving_model/penguin-pipeline2/1682589392  min value  serving_model/penguin-pipeline2/1682510985


In [171]:
# Prepare an example and run inference.

#Chinstrap test
features = {
  'culmen_length_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[46.3])),
  'culmen_depth_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[17.5])),
  'flipper_length_mm': tf.train.Feature(int64_list=tf.train.Int64List(value=[187])),
  'body_mass_g': tf.train.Feature(int64_list=tf.train.Int64List(value=[3200])),
}
example_proto = tf.train.Example(features=tf.train.Features(feature=features))
examples = example_proto.SerializeToString()

result_new = inference_fn_new(examples=tf.constant([examples]))
print('Chinstrap test: ')
print('Model with incomplete dataset result: ', result_new['output_0'].numpy())

result_old = inference_fn_old(examples=tf.constant([examples]))
print('Model with full dataset result: ', result_old['output_0'].numpy())
 
#Adelie test
features = {
  'culmen_length_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[34.9])),
  'culmen_depth_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[17.5])),
  'flipper_length_mm': tf.train.Feature(int64_list=tf.train.Int64List(value=[190])),
  'body_mass_g': tf.train.Feature(int64_list=tf.train.Int64List(value=[3723])),
}
example_proto = tf.train.Example(features=tf.train.Features(feature=features))
examples = example_proto.SerializeToString()
print('Adelie test: ')
result_new = inference_fn_new(examples=tf.constant([examples]))
print('Model with incomplete dataset result: ', result_new['output_0'].numpy())

result_old = inference_fn_old(examples=tf.constant([examples]))
print('Model with full dataset result: ', result_old['output_0'].numpy())

#Gentoo test 
features = {
  'culmen_length_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[49.5])),
  'culmen_depth_mm': tf.train.Feature(float_list=tf.train.FloatList(value=[16.5])),
  'flipper_length_mm': tf.train.Feature(int64_list=tf.train.Int64List(value=[227])),
  'body_mass_g': tf.train.Feature(int64_list=tf.train.Int64List(value=[6100])),
}
example_proto = tf.train.Example(features=tf.train.Features(feature=features))
examples = example_proto.SerializeToString()
print('Gentoo test: ')
result_new = inference_fn_new(examples=tf.constant([examples]))
print('Model with incomplete dataset result: ', result_new['output_0'].numpy())

result_old = inference_fn_old(examples=tf.constant([examples]))
print('Model with full dataset result: ', result_old['output_0'].numpy())

TypeError: signature_wrapper(*, culmen_depth_mm, body_mass_g, culmen_length_mm, flipper_length_mm) missing required arguments: body_mass_g, culmen_depth_mm, culmen_length_mm, flipper_length_mm.

In [None]:
!head{_data_filepath}

### MLMD DATABASE QUERY

In [None]:
import os
import tempfile
import urllib
import pandas as pd

import tensorflow_model_analysis as tfma
from tfx.orchestration.experimental.interactive.interactive_context import InteractiveContext


In [None]:
connection_config = tfx.orchestration.metadata.sqlite_metadata_connection_config(METADATA_PATH)
store = mlmd.MetadataStore(connection_config)

# All TFX artifacts are stored in the base directory
base_dir = connection_config.sqlite.filename_uri.split('metadata.sqlite')[0]


In [None]:
def display_types(types):
  # Helper function to render dataframes for the artifact and execution types
  table = {'id': [], 'name': []}
  for a_type in types:
    table['id'].append(a_type.id)
    table['name'].append(a_type.name)
  return pd.DataFrame(data=table)


In [None]:
def display_artifacts(store, artifacts):
  # Helper function to render dataframes for the input artifacts
  table = {'artifact id': [], 'type': [], 'uri': []}
  for a in artifacts:
    table['artifact id'].append(a.id)
    artifact_type = store.get_artifact_types_by_id([a.type_id])[0]
    table['type'].append(artifact_type.name)
    table['uri'].append(a.uri.replace(base_dir, './'))
  return pd.DataFrame(data=table)


In [None]:
def display_properties(store, node):
  # Helper function to render dataframes for artifact and execution properties
  table = {'property': [], 'value': []}
  for k, v in node.properties.items():
    table['property'].append(k)
    table['value'].append(
        v.string_value if v.HasField('string_value') else v.int_value)
  for k, v in node.custom_properties.items():
    table['property'].append(k)
    table['value'].append(
        v.string_value if v.HasField('string_value') else v.int_value)
  return pd.DataFrame(data=table)


In [None]:
display_types(store.get_artifact_types())


In [None]:
pushed_models = store.get_artifacts_by_type("PushedModel")
display_artifacts(store, pushed_models)


In [None]:
pushed_model = pushed_models[-1]
display_properties(store, pushed_model)


In [None]:
def get_one_hop_parent_artifacts(store, artifacts):
  # Get a list of artifacts within a 1-hop of the artifacts of interest
  artifact_ids = [artifact.id for artifact in artifacts]
  executions_ids = set(
      event.execution_id
      for event in store.get_events_by_artifact_ids(artifact_ids)
      if event.type == mlmd.proto.Event.OUTPUT)
  artifacts_ids = set(
      event.artifact_id
      for event in store.get_events_by_execution_ids(executions_ids)
      if event.type == mlmd.proto.Event.INPUT)
  return [artifact for artifact in store.get_artifacts_by_id(artifacts_ids)]


In [None]:
parent_artifacts = get_one_hop_parent_artifacts(store, [pushed_model])
display_artifacts(store, parent_artifacts)


In [None]:
exported_model = parent_artifacts[0]
display_properties(store, exported_model)


In [None]:
model_parents = get_one_hop_parent_artifacts(store, [exported_model])
display_artifacts(store, model_parents)


In [None]:
used_data = model_parents[0]
display_properties(store, used_data)


In [None]:
display_types(store.get_execution_types())


In [None]:
def find_producer_execution(store, artifact):
  executions_ids = set(
      event.execution_id
      for event in store.get_events_by_artifact_ids([artifact.id])
      if event.type == mlmd.proto.Event.OUTPUT)
  return store.get_executions_by_id(executions_ids)[0]

trainer = find_producer_execution(store, exported_model)
display_properties(store, trainer)
