apache_beam.ml.inference.vertex_ai_model_monitoring_v2 module

A PTransform for integrating Vertex AI Model Monitoring v2 with Apache Beam RunInference.

Vertex AI Model Monitoring v2 provides drift and skew detection on arbitrary models by evaluating input features, predictions, and attribution stats logged to BigQuery against a training baseline.

class apache_beam.ml.inference.vertex_ai_model_monitoring_v2.VertexModelMonitoringV2(project_id: str, location: str, display_name: str, model_name: str, model_version_id: str, model_monitoring_schema: Any, training_dataset: Any, tabular_objective_spec: Any, target_dataset: Any, unpack_fn: Callable[[PredictionResult], dict[str, Any]], bigquery_table: str, bigquery_schema: str | dict[str, Any] | None = None, write_to_bigquery_kwargs: dict[str, Any] | None = None, model_monitor_id: str | None = None, cron: str | None = None, schedule_display_name: str | None = None, monitoring_job_display_name: str | None = None, explanation_spec: Any | None = None, output_spec: Any | None = None, notification_spec: Any | None = None, credentials: Any | None = None, start_time: Any | None = None, end_time: Any | None = None, **kwargs)[source]

Bases: PTransform[PCollection[PredictionResult], PCollection[PredictionResult]]

A composite PTransform that exports inference outputs to BigQuery and coordinates Vertex AI Model Monitoring v2 jobs.

In batch pipelines, it blocks until inference records are committed to BigQuery before triggering an asynchronous ad-hoc monitoring job. In streaming pipelines, it provisions a recurring monitoring schedule at startup.

Parameters:
  • project_id – GCP project ID where the model monitor is created.

  • location – GCP location/region (e.g. ‘us-central1’).

  • display_name – User-visible display name for the model monitor.

  • model_name – Resource name or ID of the monitored model.

  • model_version_id – Version ID of the model.

  • model_monitoring_schema – Schema specification describing input and output features.

  • training_dataset – Baseline dataset specification (e.g. Training dataset).

  • tabular_objective_spec – Drift and skew objective parameters.

  • target_dataset – Target dataset specification pointing to production BigQuery logs.

  • unpack_fn – Callable converting PredictionResult into a dictionary matching BigQuery table schema.

  • bigquery_table – Destination BigQuery table spec in the format ‘project:dataset.table’ or ‘dataset.table’.

  • bigquery_schema – BigQuery schema definition for the destination table.

  • write_to_bigquery_kwargs – Optional dictionary of keyword arguments passed to WriteToBigQuery.

  • model_monitor_id – Optional deterministic resource ID for the model monitor. If omitted, Vertex AI generates an ID automatically.

  • cron – Cron expression defining the recurring schedule for streaming pipelines (e.g. ‘@daily’, ‘0 * * * *’). Required for streaming pipelines.

  • schedule_display_name – Display name for the streaming monitoring schedule.

  • monitoring_job_display_name – Display name for the monitoring job.

  • explanation_spec – Optional feature attribution monitoring specification.

  • output_spec – Optional output specification for monitoring statistics.

  • notification_spec – Optional alerting and notification configuration.

  • credentials – Optional google.auth credentials.

  • start_time – Optional start timestamp for streaming schedule.

  • end_time – Optional end timestamp for streaming schedule.

annotations() → dict[str, Any][source]
expand(pcoll: PCollection[PredictionResult]) → PCollection[PredictionResult][source]