Aller au contenu principal

Créer un moniteur à l'aide de l'API Databricks

Cette page explique comment créer un profilage des données dans Databricks à l'aide du Databricks SDK et décrit les paramètres utilisés dans les appels d'API. Vous pouvez également créer et gérer un profil de données à l'aide de l'API REST. Pour créer un moniteur à l'aide de l'interface utilisateur, consultez Créer un moniteur à l'aide de l'interface utilisateur Databricks.

Pour plus d'informations, consultez la référence du SDK de profilage des données et la référence de l'API REST.

Vous pouvez créer un profil sur n'importe quelle table Delta gérée ou externe enregistrée dans Unity Catalog. Un seul profil peut être créé dans un métastore Unity Catalog pour n'importe quelle table.

Exigences

Pour utiliser la version la plus récente de l'API, utilisez la commande suivante au début de votre Notebook pour installer le client Python :

Python
%pip install "databricks-sdk>=0.68.0"

Pour vous authentifier afin d'utiliser le Databricks SDK dans votre environnement, consultez la page Authentification.

Types de profil

Lorsque vous créez un profil, vous sélectionnez l'un des types de profil suivants : TimeSeries, InferenceLog ou Snapshot. Cette section décrit brièvement chaque option. Pour plus de détails, consultez la référence du SDK de profilage des données ou la référence de l’API REST.

remarque
  • Lorsque vous créez pour la première fois une série temporelle ou un profil d'inférence, Databricks analyse uniquement les données des 30 jours précédant sa création. Une fois le profil créé, toutes les nouvelles données sont traitées.
  • Les profils définis sur les vues matérialisées ne prennent pas en charge le traitement incrémentiel.
astuce

Pour les profils TimeSeries et Inference, il est recommandé d'activer le flux de données de modification (CDF) sur votre table. Lorsque le CDF est activé, seules les données nouvellement ajoutées sont traitées, au lieu de retraiter l'intégralité de la table à chaque refresh. Cela rend l'exécution plus efficace et réduit les coûts à mesure que vous montez en charge sur de nombreuses tables.

TimeSeries profil

Un profil TimeSeries compare les distributions de données sur plusieurs fenêtres temporelles. Pour un profil TimeSeries, vous devez fournir les éléments suivants :

  • Une colonne de timestamp (timestamp_column). Le type de données de la colonne Timestamp doit être TIMESTAMP ou un type qui peut être converti en Timestamp à l'aide de la to_timestamp fonction PySpark.
  • L'ensemble de granularities sur lequel calculer les métriques. Voici les granularités disponibles :
    • AGGREGATION_GRANULARITY_5_MINUTES
    • AGGREGATION_GRANULARITY_30_MINUTES
    • AGGREGATION_GRANULARITY_1_HOUR
    • GRANULARITÉ_AGRÉGATION_1_JOUR
    • AGGREGATION_GRANULARITY_1_WEEK
    • AGGREGATION_GRANULARITY_2_WEEKS
    • AGGREGATION_GRANULARITY_3_WEEKS
    • AGGREGATION_GRANULARITY_4_WEEKS
    • AGGREGATION_GRANULARITY_1_MONTH
    • GRANULARITÉ_D'AGRÉGATION_1_AN
Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.dataquality import Monitor, DataProfilingConfig, TimeSeriesConfig, AggregationGranularity, DataProfilingStatus, RefreshState, Refresh

w = WorkspaceClient()

schema = w.schemas.get(full_name=f"{catalog}.{schema}")
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")

config = DataProfilingConfig(
output_schema_id=schema.schema_id,
assets_dir=f"/Workspace/Users/{username}/databricks_quality_monitoring/{TABLE_NAME}",
time_series=TimeSeriesConfig(
timestamp_column="ts",
granularities=[AggregationGranularity.AGGREGATION_GRANULARITY_1_DAY]),
slicing_exprs=["type='Red'"]
)

info = w.data_quality.create_monitor(
monitor=Monitor(
object_type="table", # object_type is always "table" for data profiling
object_id=table.table_id,
data_profiling_config=config,
),
)

InferenceLog profil

Un profil InferenceLog est similaire à un profil TimeSeries, mais il inclut également des métriques de qualité du modèle. Les profils InferenceLog utilisent les paramètres suivants :

parameter

Description

problem_type

MonitorInferenceLogProblemType.PROBLEM_TYPE_CLASSIFICATION ou MonitorInferenceLogProblemType.PROBLEM_TYPE_REGRESSION

prediction_column

Colonne contenant les valeurs prédites du modèle.

timestamp_column

Colonne contenant le Timestamp de la demande d'inférence.

model_id_column

Colonne contenant l'ID du modèle utilisé pour la prédiction.

granularities

Détermine comment partitionner les données dans des fenêtres temporelles. Consultez le TimeSeries profil pour connaître les valeurs disponibles.

label_column

(Facultatif) Colonne contenant la vérité terrain pour les prédictions du modèle.

parameter

Description

problem_type

MonitorInferenceLogProblemType.PROBLEM_TYPE_CLASSIFICATION ou MonitorInferenceLogProblemType.PROBLEM_TYPE_REGRESSION

prediction_column

Colonne contenant les valeurs prédites du modèle.

timestamp_column

Colonne contenant le Timestamp de la demande d'inférence.

model_id_column

Colonne contenant l'ID du modèle utilisé pour la prédiction.

granularities

Détermine comment partitionner les données dans des fenêtres temporelles. Consultez le TimeSeries profil pour connaître les valeurs disponibles.

label_column

(Facultatif) Colonne contenant la vérité terrain pour les prédictions du modèle.

Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.dataquality import Monitor, DataProfilingConfig, InferenceLogConfig, InferenceProblemType, AggregationGranularity, DataProfilingStatus, RefreshState, Refresh

w = WorkspaceClient()

schema = w.schemas.get(full_name=f"{catalog}.{schema}")
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")

config = DataProfilingConfig(
output_schema_id=schema.schema_id,
assets_dir=f"/Workspace/Users/{username}/databricks_quality_monitoring/{TABLE_NAME}",
inference_log=InferenceLogConfig(
problem_type=InferenceProblemType.INFERENCE_PROBLEM_TYPE_CLASSIFICATION,
prediction_column="preds",
model_id_column="model_ver",
label_column="label", # optional
timestamp_column="ts",
granularities=[AggregationGranularity.AGGREGATION_GRANULARITY_1_DAY])
)

info = w.data_quality.create_monitor(
monitor=Monitor(
object_type="table",
object_id=table.table_id,
data_profiling_config=config,
),
)

Pour InferenceLog profils, les tranches sont automatiquement créées en fonction des valeurs distinctes de model_id_col.

Snapshot profil

Contrairement à TimeSeries, un Snapshot profile comment l'intégralité du contenu de la table évolue au fil du temps. Les métriques sont calculées sur toutes les données de la table et reflètent l'état de la table à chaque actualisation du profil.

remarque

La taille maximale de table pour un profil d'instantané est de 4 To. Pour les tables plus volumineuses, utilisez plutôt des profils de séries temporelles.

Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.dataquality import Monitor, DataProfilingConfig, SnapshotConfig, DataProfilingStatus, RefreshState, Refresh

w = WorkspaceClient()

schema = w.schemas.get(full_name=f"{catalog}.{schema}")
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")
table_id = table.table_id
table_object_type = "table"

config = DataProfilingConfig(
output_schema_id=schema.schema_id,
assets_dir=f"/Workspace/Users/{username}/databricks_quality_monitoring/{TABLE_NAME}",
snapshot=SnapshotConfig(),
slicing_exprs=["type='Red'"]
)

refresh et afficher les résultats

Pour consulter l'historique de refresh, vous devez utiliser le workspace Databricks à partir duquel le profilage des données a été activé.

Pour refresh les tables de métriques, utilisez create_refresh. Par exemple :

Python
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
run_info = w.data_quality.create_refresh(
object_type=table_object_type, object_id=table_id, refresh=Refresh(
object_type=table_object_type,
object_id=table_id,
)
)

Lorsque vous appelez create_refresh à partir d'un Notebook, les tables de métriques sont créées ou mises à jour. Ce calcul s'exécute sur le compute Serverless, et non sur le cluster auquel le Notebook est attaché. Vous pouvez continuer à exécuter des commandes dans le Notebook pendant que les statistiques sont mises à jour.

Les tables de métriques sont des tables Unity Catalog. Vous pouvez les query dans des Notebooks ou dans l'explorateur de query SQL, et les afficher dans Catalog Explorer.

Pour afficher l'historique de tous les refresh associés à un profil, utilisez list_refreshes.

Python
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
it = w.data_quality.list_refresh(object_type=table_object_type, object_id=table_id)

Pour obtenir le statut d'une exécution spécifique qui a été mise en file d'attente, en cours d'exécution ou terminée, utilisez get_refresh.

Python
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
it = w.data_quality.list_refresh(object_type=table_object_type, object_id=table_id)

run_info = next(it, None)
while run_info.state in (RefreshState.MONITOR_REFRESH_STATE_PENDING, RefreshState.MONITOR_REFRESH_STATE_RUNNING):
run_info = w.data_quality.get_refresh(object_type=table_object_type, object_id=table_id, refresh_id=run_info.refresh_id)
time.sleep(30)

Afficher les paramètres de profil

Vous pouvez examiner les paramètres de profil à l'aide de l'API get_monitor.

Python
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")
w.data_quality.get_monitor(object_type="table", object_id=table.table_id)

Planifier

Pour configurer un profil à exécuter de manière planifiée, utilisez le parameter schedule de create_monitor:

Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.catalog import MonitorTimeSeries, MonitorCronSchedule

w = WorkspaceClient()
schema = w.schemas.get(full_name=f"{catalog}.{schema}")
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")

config = DataProfilingConfig(
output_schema_id=schema.schema_id,
snapshot=SnapshotConfig(),
schedule=CronSchedule(
quartz_cron_expression="0 0 12 * * ?", # schedules a refresh every day at 12 noon
timezone_id="PST",
)
)

info = w.data_quality.create_monitor(
monitor=Monitor(
object_type="table",
object_id=table.table_id,
data_profiling_config=config,
),
)

Consultez les expressions cron pour plus d'informations.

Notifications

Pour configurer les notifications d'un profil, utilisez le paramètre notifications de create_monitor:

Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.dataquality import Monitor, DataProfilingConfig, SnapshotConfig, NotificationSettings, NotificationDestination

w = WorkspaceClient()
schema = w.schemas.get(full_name=f"{catalog}.{schema}")
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")

config = DataProfilingConfig(
output_schema_id=schema.schema_id,
snapshot=SnapshotConfig(),
notification_settings=NotificationSettings(
# Notify the given email when a monitoring refresh fails or times out.
on_failure=NotificationDestination(
email_addresses=["your_email@domain.com"]
)
)
)

info = w.data_quality.create_monitor(
monitor=Monitor(
object_type="table",
object_id=table.table_id,
data_profiling_config=config,
),
)

Un maximum de 5 adresses e-mail est pris en charge par type d'événement (par exemple, « on_failure »).

Contrôler l'accès aux tables métriques

Les tables de métriques et le tableau de bord créés par un profil sont la propriété de l'utilisateur qui a créé le profil. Vous pouvez utiliser les privilèges Unity Catalog pour contrôler l'accès aux tables de métriques. Pour partager des tableaux de bord au sein d'un Workspace, utilisez le bouton Partager en haut à droite du tableau de bord.

Supprimer un profil

Pour supprimer un profil :

Python
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
table = w.tables.get(full_name=f"{catalog}.{schema}.{table_name}")
w.data_quality.delete_monitor(object_type="table", object_id=table.table_id)

Cette commande ne supprime pas les tables de profil et le tableau de bord créés par le profil. Vous devez supprimer ces assets lors d'une étape distincte, ou vous pouvez les enregistrer à un autre emplacement.