Aller au contenu principal

Surveiller la progression de la passerelle d'ingestion avec les logs d'événements

S’applique à : Icône X rouge connecteurs SaaS Icône de coche verte connecteurs de base de données

Découvrez comment utiliser les Logs d'événements pour monitorer la progression des passerelles d'ingestion en temps réel. Les Logs d'événements fournissent des métriques par table pour les phases d'instantané et de capture de données modifiées (CDC), vous permettant de suivre la santé du pipeline, d'identifier les pipelines bloqués et de créer des Solutions de monitoring automatisées.

Les événements de progression vous permettent de :

  • Suivez le nombre de lignes et d'octets ingérés par table sans attendre l'achèvement du pipeline.
  • Surveillez la progression de l'instantané pour chaque table afin d'estimer l'achèvement des gros chargements initiaux.
  • Estimez quand un instantané de longue durée sera terminé en utilisant l'ETA par table.
  • Mesurer la latence de découverte CDC de bout en bout (du commit source à l'émission des logs d'événements) par table.
  • Surveillez chaque table ingérée individuellement pour identifier les goulots d'affichage ou les problèmes.
  • Recevez des événements même si aucune modification de données ne se produit pour confirmer que le pipeline est en cours d'exécution.
  • Créez des alertes et des tableaux de bord à l'aide de données d'événements structurées au lieu d'analyser les logs.

Fonctionnement des événements de progression

La passerelle émet les types d'événements suivants à intervalles réguliers (default : 5 minutes) pour chaque table de votre pipeline :

  • flow_progress Les événements ** ** signalent les compteurs de lignes et d'octets pour les flux d'instantanés et de CDC. Les métriques de ces événements sont des deltas. Ils Reset à zéro après chaque émission. Pour les flux CDC, ces événements incluent également des métriques de latence qui mesurent la latence de découverte de bout en bout et les performances du pipeline d'upload.
  • operation_progress ** ** événements signalent la progression de l'instantané en pourcentage. Les flux d'instantanés émettent ces événements en plus de flow_progress. Le pourcentage de progression est cumulatif. Il s'accumule de 0 à 100 pendant la durée de vie de l'instantané. Ces événements incluent également le temps restant estimé (estimated_completion_ms) jusqu'à la fin de l'instantané.

Chaque événement comprend :

  • Noms de table source et de destination.
  • Métriques par table : lignes fusionnées, lignes supprimées (CDC uniquement), octets de sortie et pourcentage de progression (pour l'instantané).
  • Pour les flux CDC, les métriques de latence incluent la latence de découverte et le temps de traitement batch.
  • Pour les instantanés, l'estimation du temps restant jusqu'à l'achèvement.
  • Lorsque l'événement a été généré.

Les événements sont disponibles dans la table du journal des événements, mais pas via les APIs publiques. Vous pouvez interroger la table de log des événements en utilisant SQL pour analyser le comportement du pipeline et construire des solutions de monitoring.

Accéder aux événements de progression

Les événements de progression sont stockés dans la table des Logs d'événements. Pour y accéder :

  1. Accédez à votre passerelle dans le Databricks Workspace.
  2. Cliquez sur l'onglet Event Logs pour afficher les événements dans l'interface utilisateur.
  3. Query la table des logs d'événements directement à l'aide de SQL pour une analyse détaillée.

Interroger la table du log des événements

Pour interroger flow_progress événements pour les compteurs de lignes et d'octets :

SQL
SELECT
timestamp,
CONCAT(origin.catalog_name, '.', origin.schema_name, '.', origin.dataset_name) AS table_name,
details:flow_progress:metrics:num_upserted_rows::bigint AS rows_upserted,
COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0) AS rows_deleted,
details:flow_progress:metrics:num_output_bytes::bigint AS output_bytes,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'snapshot'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'cdc'
ELSE 'unknown'
END AS ingestion_phase
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.pipeline_type = 'INGESTION_GATEWAY'
ORDER BY timestamp DESC

To query operation_progress events for snapshot progress percentage:

SQL
SELECT
timestamp,
origin.flow_name AS flow_name,
details:operation_progress:status::string AS status,
details:operation_progress:progress_percent::double AS progress_pct
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.pipeline_type = 'INGESTION_GATEWAY'
ORDER BY timestamp DESC

Remplacez <pipeline-id> par l'ID de votre passerelle.

Comprendre la structure d'événement

Les événements de progression utilisent l'un des types d'événements suivants avec le niveau de Logs METRICS :

  • flow_progress: émis pour les flux de snapshot et de CDC. Indique les deltas de lignes et d'octets par table.
  • operation_progress: Émis uniquement pour les flux de snapshot. Indique le pourcentage d'achèvement de l'instantané pour une table.

Les exemples suivants montrent la structure JSON pour chaque type d'événement :

Structure d'événement de progression de flux d'instantané

JSON
{
"id": "01234567-89ab-cdef-0123-456789abcdef",
"timestamp": "2025-10-14T13:33:14.175Z",
"level": "METRICS",
"event_type": "flow_progress",
"origin": {
"pipeline_type": "INGESTION_GATEWAY",
"pipeline_name": "MyPipeline",
"dataset_name": "customers",
"catalog_name": "main",
"schema_name": "sales",
"flow_name": "main.sales.customers_snapshot_flow",
"ingestion_source_type": "SQLSERVER"
},
"message": "Completed a streaming update of 'main.sales.customers_snapshot_flow'.",
"details": {
"flow_progress": {
"status": "RUNNING",
"metrics": {
"num_upserted_rows": 7512704,
"num_deleted_rows": null,
"num_output_bytes": 458752000
}
}
},
"maturity_level": "STABLE"
}

Structure d'événement de progression du flux CDC

JSON
{
"id": "01234567-89ab-cdef-0123-456789abcdef",
"timestamp": "2025-10-14T13:33:57.426Z",
"level": "METRICS",
"event_type": "flow_progress",
"origin": {
"pipeline_type": "INGESTION_GATEWAY",
"pipeline_name": "MyPipeline",
"dataset_name": "customers",
"catalog_name": "main",
"schema_name": "sales",
"flow_name": "main.sales.customers_cdc_flow",
"ingestion_source_type": "SQLSERVER"
},
"message": "Completed a streaming update of 'main.sales.customers_cdc_flow'.",
"details": {
"flow_progress": {
"status": "RUNNING",
"metrics": {
"num_upserted_rows": 25,
"num_deleted_rows": 3,
"num_output_bytes": 18432
},
"streaming_metrics": {
"discovery_latency_ms": 12450,
"batch_processing_time_ms": 8100,
"event_time": {
"max": "2025-10-14T13:33:45.000Z"
}
}
}
},
"maturity_level": "STABLE"
}

Structure d'événement de progression d'Opération Snapshot

JSON
{
"id": "01234567-89ab-cdef-0123-456789abcdef",
"timestamp": "2025-10-14T13:33:14.175Z",
"level": "METRICS",
"event_type": "operation_progress",
"origin": {
"pipeline_type": "INGESTION_GATEWAY",
"pipeline_name": "MyPipeline",
"dataset_name": "customers",
"catalog_name": "main",
"schema_name": "sales",
"flow_name": "main.sales.customers_snapshot_flow",
"ingestion_source_type": "SQLSERVER"
},
"message": "Snapshot in progress for 'main.sales.customers'.",
"details": {
"operation_progress": {
"type": "CDC_SNAPSHOT",
"status": "IN_PROGRESS",
"duration_ms": 3600000,
"progress_percent": 65.5,
"estimated_completion_ms": 1885000,
"cdc_snapshot": {
"target_table_name": "main.sales.customers",
"snapshot_timestamp": 1737542400000,
"snapshot_reason": "NEW_TABLE"
}
}
},
"maturity_level": "STABLE"
}

Champs d'événement

Le tableau suivant décrit les champs clés des événements en cours :

Champ

Type

Description

event_type

Chaîne

Soit flow_progress (compteurs de lignes et d'octets pour l'instantané et la CDC), soit operation_progress (pourcentage de progression de l'instantané).

level

Chaîne

Toujours METRICS.

timestamp

Chaîne

Timestamp ISO 8601 lorsque l'événement a été généré.

origin.pipeline_type

Chaîne

Toujours INGESTION_GATEWAY.

origin.pipeline_name

Chaîne

Nom de la passerelle.

origin.dataset_name

Chaîne

Nom de la table ingérée.

origin.catalog_name

Chaîne

Nom du catalogue Unity Catalog.

origin.schema_name

Chaîne

Nom du schéma Unity Catalog.

origin.flow_name

Chaîne

Identifiant de flux qui indique la phase d'ingestion. Format : {catalog}.{schema}.{table}_snapshot_flow pour le chargement initial ou {catalog}.{schema}.{table}_cdc_flow pour les modifications incrémentielles.

origin.ingestion_source_type

Chaîne

Type de base de données source (par exemple, SQLSERVER, MYSQL, POSTGRESQL, ORACLE).

details:flow_progress.status

Chaîne

Statut de flux actuel, généralement RUNNING.

details:flow_progress.metrics.num_upserted_rows

Entier

Nombre de lignes insérées ou mises à jour depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission.

details:flow_progress.metrics.num_deleted_rows

Entier

Nombre de lignes supprimées depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission. null pour les flux d'instantanés (l'instantané ne supprime pas).

details:flow_progress.metrics.num_output_bytes

Entier

Nombre d'octets compressés upload vers un volume depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission. Rempli pour les flux de snapshot et de CDC.

details:flow_progress.streaming_metrics.discovery_latency_ms

Entier

Temps en millisecondes entre le changement de source à event_time.max et le moment où cet événement a été émis. Flux CDC uniquement. Peut être null si la source n'expose pas les horodatages de modification.

details:flow_progress.streaming_metrics.batch_processing_time_ms

Entier

Temps en millisecondes que la passerelle a passé à lire et à upload le dernier batch de modifications. N'inclut pas le délai côté source. Flux CDC uniquement. Peut être null jusqu'à ce que le premier batch soit terminé.

details:flow_progress.streaming_metrics.event_time.max

Chaîne

Timestamp ISO 8601 de la modification de source la plus récente que la passerelle a lue pour cette table. Flux CDC seulement.

details:operation_progress.type

Chaîne

Type d'opération. CDC_SNAPSHOT pour les opérations d'instantané.

details:operation_progress.status

Chaîne

Statut actuel de l'opération. IN_PROGRESS pendant que l'instantané est en cours d'exécution, COMPLETED lorsqu'il se termine. D'autres valeurs incluent STARTED, CANCELED et FAILED.

details:operation_progress.duration_ms

Entier

Temps total écoulé de l'opération en millisecondes.

details:operation_progress.progress_percent

Double

Pourcentage d'achèvement de l'instantané (0.0 - 100.0). Cumulatif, pas un delta. La valeur augmente à mesure que les segments sont achevés et atteint 100.0 lorsque l'instantané est terminé.

details:operation_progress.estimated_completion_ms

Entier

Temps estimé restant en millisecondes jusqu'à la fin de l'instantané. Diminue au fur et à mesure que l'instantané progresse et atteint 0 à la fin. Peut augmenter brièvement si les progrès ralentissent. Peut être null jusqu'à ce que suffisamment de données aient été traitées pour produire une estimation.

details:operation_progress.cdc_snapshot.target_table_name

Chaîne

Nom entièrement qualifié de la table faisant l'objet d'une prise d'instantané.

maturity_level

Chaîne

Toujours STABLE.

Champ

Type

Description

event_type

Chaîne

Soit flow_progress (compteurs de lignes et d'octets pour l'instantané et la CDC), soit operation_progress (pourcentage de progression de l'instantané).

level

Chaîne

Toujours METRICS.

timestamp

Chaîne

Timestamp ISO 8601 lorsque l'événement a été généré.

origin.pipeline_type

Chaîne

Toujours INGESTION_GATEWAY.

origin.pipeline_name

Chaîne

Nom de la passerelle.

origin.dataset_name

Chaîne

Nom de la table ingérée.

origin.catalog_name

Chaîne

Nom du catalogue Unity Catalog.

origin.schema_name

Chaîne

Nom du schéma Unity Catalog.

origin.flow_name

Chaîne

Identifiant de flux qui indique la phase d'ingestion. Format : {catalog}.{schema}.{table}_snapshot_flow pour le chargement initial ou {catalog}.{schema}.{table}_cdc_flow pour les modifications incrémentielles.

origin.ingestion_source_type

Chaîne

Type de base de données source (par exemple, SQLSERVER, MYSQL, POSTGRESQL, ORACLE).

details:flow_progress.status

Chaîne

Statut de flux actuel, généralement RUNNING.

details:flow_progress.metrics.num_upserted_rows

Entier

Nombre de lignes insérées ou mises à jour depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission.

details:flow_progress.metrics.num_deleted_rows

Entier

Nombre de lignes supprimées depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission. null pour les flux d'instantanés (l'instantané ne supprime pas).

details:flow_progress.metrics.num_output_bytes

Entier

Nombre d'octets compressés upload vers un volume depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission. Rempli pour les flux de snapshot et de CDC.

details:flow_progress.streaming_metrics.discovery_latency_ms

Entier

Temps en millisecondes entre le changement de source à event_time.max et le moment où cet événement a été émis. Flux CDC uniquement. Peut être null si la source n'expose pas les horodatages de modification.

details:flow_progress.streaming_metrics.batch_processing_time_ms

Entier

Temps en millisecondes que la passerelle a passé à lire et à upload le dernier batch de modifications. N'inclut pas le délai côté source. Flux CDC uniquement. Peut être null jusqu'à ce que le premier batch soit terminé.

details:flow_progress.streaming_metrics.event_time.max

Chaîne

Timestamp ISO 8601 de la modification de source la plus récente que la passerelle a lue pour cette table. Flux CDC seulement.

details:operation_progress.type

Chaîne

Type d'opération. CDC_SNAPSHOT pour les opérations d'instantané.

details:operation_progress.status

Chaîne

Statut actuel de l'opération. IN_PROGRESS pendant que l'instantané est en cours d'exécution, COMPLETED lorsqu'il se termine. D'autres valeurs incluent STARTED, CANCELED et FAILED.

details:operation_progress.duration_ms

Entier

Temps total écoulé de l'opération en millisecondes.

details:operation_progress.progress_percent

Double

Pourcentage d'achèvement de l'instantané (0.0 - 100.0). Cumulatif, pas un delta. La valeur augmente à mesure que les segments sont achevés et atteint 100.0 lorsque l'instantané est terminé.

details:operation_progress.estimated_completion_ms

Entier

Temps estimé restant en millisecondes jusqu'à la fin de l'instantané. Diminue au fur et à mesure que l'instantané progresse et atteint 0 à la fin. Peut augmenter brièvement si les progrès ralentissent. Peut être null jusqu'à ce que suffisamment de données aient été traitées pour produire une estimation.

details:operation_progress.cdc_snapshot.target_table_name

Chaîne

Nom entièrement qualifié de la table faisant l'objet d'une prise d'instantané.

maturity_level

Chaîne

Toujours STABLE.

Comportement des métriques

Les métriques de progression entrent dans les catégories suivantes :

Mesures Delta (num_upserted_rows, num_deleted_rows, num_output_bytes) :

  • Représenter les changements depuis le dernier événement, et non les totaux cumulatifs.
  • Reset à zéro après chaque émission d'événement.
  • Sont émises même lorsque aucune modification de données ne se produit, servant d'indicateurs de vivacité.
  • Pour les flux d'instantanés, num_deleted_rows est null car l'instantané ne produit pas de suppressions.

Métriques cumulatives (progress_percent) :

  • La valeur s'accumule de 0.0 à 100.0 pendant la durée de vie d'un instantané.
  • Mises à jour à mesure que l'instantané progresse. Pour les petites tables, la valeur peut passer directement de 0.0 à 100.0. Pour les grandes tables, la valeur se met à jour progressivement et constitue une approximation, non un nombre de lignes exact.

Mesures à un instant précis (discovery_latency_ms, batch_processing_time_ms, event_time.max, estimated_completion_ms) :

  • Chaque valeur reflète l'état de la métrique au moment où l'événement a été émis.
  • discovery_latency_ms et batch_processing_time_ms s'appliquent uniquement aux flux CDC. Les valeurs sont rapportées comme 0 si le résultat était autrement négatif.
  • event_time.max S'applique uniquement aux flux CDC. La valeur est le Timestamp de la modification de source la plus récente que la passerelle a lue.
  • estimated_completion_ms s'applique uniquement aux flux d'instantanés. La valeur diminue à mesure que l'instantané progresse et atteint 0 à la fin.

Configurer les événements de progression

Les événements de progression sont activés par default pour les nouvelles passerelles. Vous pouvez personnaliser le comportement des événements à l'aide des paramètres de configuration du pipeline.

Activer ou désactiver les événements de progression

JSON
"configuration": {
"pipelines.gateway.progressEventsEnabled": "true"
}

Définissez sur "false" pour désactiver les événements de progression.

Ajuster la fréquence d'émission des événements

JSON
"configuration": {
"pipelines.gateway.progressEventEmitFrequencySeconds": "300"
}

default : 300 secondes (cinq minutes). Plage valide : de 30 à 3 600 secondes (de 30 secondes à une heure). Ce paramètre contrôle la cadence des événements flow_progress et operation_progress.

Exemple de configuration de passerelle

L'exemple suivant montre une configuration complète de passerelle avec des événements de progression activés et définis pour s'émettre toutes les cinq minutes :

Python
gateway_pipeline_spec = {
"pipeline_type": "INGESTION_GATEWAY",
"name": "my_gateway_pipeline",
"catalog": "main",
"target": "my_schema",
"continuous": True,
"configuration": {
"pipelines.gateway.progressEventsEnabled": "true",
"pipelines.gateway.progressEventEmitFrequencySeconds": "300"
},
# ... rest of pipeline spec
}

Comportements et limitations importants

default behavior

  • La fonctionnalité est activée par default pour toutes les nouvelles passerelles.
  • Les pipelines existants reçoivent automatiquement cette fonctionnalité lors de leur prochaine mise à jour ou de leur prochain redémarrage.
  • Aucune action n’est requise pour activer les événements de progression.

Disponibilité des métriques sur les différentes versions de la passerelle

L'ETA des instantanés (estimated_completion_ms) et les métriques de latence CDC (streaming_metrics) nécessitent une image de passerelle à partir de la version de la passerelle d'ingestion de mai 2026 ou ultérieure. Databricks sélectionne automatiquement l'image de la passerelle. Il est impossible de configurer cela manuellement. Les fonctionnalités ont été déployées dans toutes les régions de production en mai 2026.

Les pipelines nouvellement créés et existants reçoivent la nouvelle image. Les pipelines existants l'adoptent automatiquement lors de leur prochaine mise à jour ou de leur prochain redémarrage, vous n'avez donc pas besoin de migrer.

Pour vérifier si votre pipeline émet les champs estimated_completion_ms et streaming_metrics, exécutez la query suivante. Si les deux colonnes retournent des lignes, votre passerelle prend en charge les métriques d'ETA des instantanés et de latence CDC :

SQL
SELECT
MAX(CASE WHEN details:operation_progress:estimated_completion_ms IS NOT NULL
THEN timestamp END) AS last_snapshot_eta,
MAX(CASE WHEN details:flow_progress:streaming_metrics IS NOT NULL
THEN timestamp END) AS last_cdc_latency
FROM event_log('<pipeline-id>')
WHERE event_type IN ('operation_progress', 'flow_progress')
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 1 HOUR

Si aucune colonne n’a de Timestamp récent après un intervalle d’émission complet (default : cinq minutes) suite à un redémarrage de pipeline, contactez le support Databricks pour confirmer si les fonctionnalités sont activées dans votre région.

Considérations relatives au timing

  • La première émission peut prendre jusqu'à l'intervalle de fréquence configuré (default : cinq minutes) après le start du pipeline avant que les événements de progression n'apparaissent.
  • Les événements sont émis à la fréquence configurée pendant l'ingestion active.

Métriques sans mise à jour

  • Les événements sont émis pour toutes les tables, y compris celles avec zéro mise à jour.

  • Les métriques de non-mise à jour aident à faire la distinction entre :

    • Tables inactives : traitées, mais aucune modification de données n'est survenue.
    • Tables non traitées : Non encore prises en charge par le pipeline.
  • Les événements de non-mise à jour servent de signaux de maintien en vie confirmant que le pipeline est en cours d'exécution active.

Comportement du pourcentage de progression du snapshot

  • La progression de l’instantané est calculée comme (completed_chunks / total_chunks) × 100. La métrique est approximative, pas un pourcentage exact au niveau des lignes.
  • Les tables qui ne sont pas divisées en plusieurs segments (généralement des tables plus petites) passent directement de 0.0 à 100.0 entre les émissions, car il n'y a qu'un seul segment à suivre.
  • Les grandes tables divisées en plusieurs blocs se mettent à jour de manière incrémentielle à mesure que chaque bloc se termine, offrant un signal de progression graduel utile pour le monitoring des chargements initiaux de longue durée.
  • Un statut COMPLETED signale toujours progress_percent = 100.0.
  • La métrique ne survit pas à un refresh ou un redémarrage du pipeline. Après un redémarrage, la progression de l'instantané reprend à partir du dernier point de contrôle validé et la métrique continue de grimper à partir de la position reprise.

Exemples de query

Les requêtes d'exemple suivantes montrent comment surveiller votre passerelle à l'aide du nombre de lignes, des octets de sortie, de la progression et de l'ETA des instantanés, ainsi que des métriques de latence de la CDC. Remplacez <pipeline-id> par votre ID de passerelle dans chaque query.

astuce

Le Notebook Moniteur de progression de la passerelle d'ingestion contient toutes les requêtes de cette section. Importez le Notebook dans le workspace où votre passerelle s'exécute et spécifiez votre ID de passerelle. Exécutez le Notebook pour inspecter la progression de l'instantané, la latence de la CDC et le nombre de lignes. Vous pouvez également ajuster les thresholds de SLA pour correspondre à vos exigences de monitoring.

Queries pour le nombre de lignes

Volume par table (dernières 24 heures)

Total des opérations d'upsert et de suppression pour chaque table au cours des dernières 24 heures, la phase d'ingestion étant classée à partir du suffixe du nom du flux. Utilisez ceci comme titre de tableau de bord pour voir quelles tables ont déplacé le plus de données.

SQL
WITH row_events AS (
SELECT
origin.flow_name AS flow_name,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'SNAPSHOT'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'CDC'
ELSE 'OTHER'
END AS phase,
details:flow_progress:metrics:num_upserted_rows::bigint AS upserts,
details:flow_progress:metrics:num_deleted_rows::bigint AS deletes
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
)
SELECT
flow_name,
phase,
SUM(upserts) AS rows_upserted_24h,
SUM(COALESCE(deletes, 0)) AS rows_deleted_24h,
SUM(upserts) + SUM(COALESCE(deletes, 0)) AS total_rows_moved_24h
FROM row_events
GROUP BY flow_name, phase
ORDER BY total_rows_moved_24h DESC

Événements de progression récents (dernière heure)

Compteurs de lignes récents pour toutes les tables de votre pipeline. Utile pour le monitoring en quasi temps réel.

SQL
SELECT
origin.pipeline_name,
origin.dataset_name,
origin.flow_name,
details:flow_progress:metrics:num_upserted_rows::bigint AS num_upserted_rows,
COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0) AS num_deleted_rows,
timestamp
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND timestamp >= current_timestamp() - INTERVAL 1 HOUR
ORDER BY timestamp DESC

Identifier les tables silencieuses ou bloquées

Les tables émettant des événements, mais ne signalant aucune mise à jour ni aucune suppression depuis les 60 dernières minutes, sont candidates au statut « bloqué ». Vérifiez si la source doit être modifiée. Les tables CDC qui sont réellement inactives (par exemple, pendant la nuit) apparaissent également ici.

SQL
WITH recent_window AS (
SELECT
origin.flow_name AS flow_name,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'SNAPSHOT'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'CDC'
ELSE 'OTHER'
END AS phase,
COUNT(*) AS emissions_in_window,
SUM(details:flow_progress:metrics:num_upserted_rows::bigint) AS upserts_in_window,
SUM(COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0)) AS deletes_in_window,
MAX(timestamp) AS last_event
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 60 MINUTES
GROUP BY origin.flow_name
)
SELECT
flow_name,
phase,
emissions_in_window,
upserts_in_window,
deletes_in_window,
last_event,
ROUND(TIMESTAMPDIFF(MINUTE, last_event, current_timestamp()), 0) AS minutes_since_last_event
FROM recent_window
WHERE upserts_in_window = 0
AND deletes_in_window = 0
ORDER BY minutes_since_last_event DESC

Chronologie par table avec totaux cumulés

Historique complet événement par événement pour un seul flux au cours des dernières 24 heures, avec des totaux de lignes cumulés. Remplacez <flow-pattern> par un modèle SQL LIKE (par exemple, '%customers%_cdc_flow').

SQL
SELECT
origin.flow_name AS flow_name,
origin.update_id AS update_id,
timestamp,
details:flow_progress:metrics:num_upserted_rows::bigint AS upserts_this_period,
COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0) AS deletes_this_period,
SUM(details:flow_progress:metrics:num_upserted_rows::bigint)
OVER (PARTITION BY origin.flow_name, origin.update_id
ORDER BY timestamp
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cumulative_upserts_this_run,
SUM(COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0))
OVER (PARTITION BY origin.flow_name, origin.update_id
ORDER BY timestamp
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cumulative_deletes_this_run
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.flow_name LIKE '<flow-pattern>'
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
ORDER BY timestamp

Requêtes en octets de sortie

Octets par table (dernières 24 heures)

Total d'octets transférés vers un volume par table au cours des dernières 24 heures, avec des unités MB et GB conviviales. Trier par ordre décroissant pour voir les tables au volume le plus élevé.

SQL
WITH byte_events AS (
SELECT
origin.flow_name AS flow_name,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'SNAPSHOT'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'CDC'
ELSE 'OTHER'
END AS phase,
details:flow_progress:metrics:num_output_bytes::bigint AS output_bytes
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
)
SELECT
flow_name,
phase,
SUM(output_bytes) AS bytes_24h,
ROUND(SUM(output_bytes) / 1024.0 / 1024.0, 2) AS mb_24h,
ROUND(SUM(output_bytes) / 1024.0 / 1024.0 / 1024.0, 3) AS gb_24h
FROM byte_events
GROUP BY flow_name, phase
ORDER BY bytes_24h DESC

Tendance du throughput (Mo par minute)

Série chronologique par minute des octets upload sur tous les flux au cours des dernières 24 heures. Rendu sous forme de graphique linéaire pour repérer les modèles de throughput et les blocages.

SQL
SELECT
DATE_TRUNC('MINUTE', timestamp) AS ts_minute,
origin.flow_name AS flow_name,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'SNAPSHOT'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'CDC'
ELSE 'OTHER'
END AS phase,
ROUND(SUM(details:flow_progress:metrics:num_output_bytes::bigint) / 1024.0 / 1024.0, 2) AS mb_per_minute
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY DATE_TRUNC('MINUTE', timestamp), origin.flow_name
ORDER BY origin.flow_name, ts_minute

Octets moyens par ligne (table large ou détecteur LOB)

Combine num_output_bytes avec les nombres de lignes pour compute la moyenne d'octets par ligne et par table. Les valeurs élevées indiquent généralement des tables LOB ou à schéma large qui influent sur le coût du throughput. Utile pour la planification de la capacité et la révision du schéma.

SQL
WITH joined AS (
SELECT
origin.flow_name AS flow_name,
CASE
WHEN origin.flow_name LIKE '%_snapshot_flow' THEN 'SNAPSHOT'
WHEN origin.flow_name LIKE '%_cdc_flow' THEN 'CDC'
ELSE 'OTHER'
END AS phase,
SUM(details:flow_progress:metrics:num_output_bytes::bigint) AS total_bytes,
SUM(details:flow_progress:metrics:num_upserted_rows::bigint) AS total_upserts,
SUM(COALESCE(details:flow_progress:metrics:num_deleted_rows::bigint, 0)) AS total_deletes
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY origin.flow_name
)
SELECT
flow_name,
phase,
total_upserts + total_deletes AS total_rows,
ROUND(total_bytes / 1024.0 / 1024.0, 2) AS total_mb,
ROUND(total_bytes / NULLIF(total_upserts + total_deletes, 0), 0) AS avg_bytes_per_row
FROM joined
WHERE total_bytes > 0
ORDER BY avg_bytes_per_row DESC

Instantanés des query de progression

Progression globale de l'instantané

La query suivante renvoie un résumé sur une seule ligne du nombre de tables terminées, en cours ou en file d'attente. Les résultats reflètent le dernier état signalé par table sur toutes les mises à jour de pipeline dans la fenêtre de rétention des logs d'événements. Ainsi, les tables terminées lors d'une mise à jour antérieure continuent de compter pour tables_completed après un refresh ou un redémarrage.

SQL
WITH latest_per_table AS (
SELECT
origin.flow_name AS flow_name,
details:operation_progress:status::string AS status,
details:operation_progress:progress_percent::double AS progress_pct,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
)
SELECT
COUNT(*) AS total_tables,
SUM(CASE WHEN status = 'COMPLETED' THEN 1 ELSE 0 END) AS tables_completed,
SUM(CASE WHEN status = 'IN_PROGRESS' AND progress_pct > 0 AND progress_pct < 100 THEN 1 ELSE 0 END) AS tables_in_progress,
SUM(CASE WHEN progress_pct = 0 THEN 1 ELSE 0 END) AS tables_not_started,
ROUND(AVG(progress_pct), 2) AS overall_progress_pct,
ROUND(SUM(CASE WHEN status = 'COMPLETED' THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 1) AS pct_tables_done
FROM latest_per_table
WHERE rn = 1

Tableau d'état des instantanés par table

La query suivante renvoie chaque table du pipeline avec son dernier statut signalé et son pourcentage de progression. Les résultats incluent les tables complétées lors des mises à jour précédentes, de sorte qu'aucune table n'apparaît comme des lignes manquantes.

SQL
WITH table_status AS (
SELECT
origin.flow_name AS flow_name,
details:operation_progress:status::string AS status,
details:operation_progress:progress_percent::double AS progress_pct,
timestamp,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
)
SELECT
flow_name,
status,
ROUND(progress_pct, 2) AS progress_pct,
timestamp AS last_update,
CASE
WHEN status = 'COMPLETED' THEN 'Done'
WHEN progress_pct = 0 THEN 'Queued'
ELSE 'Active'
END AS phase
FROM table_status
WHERE rn = 1
ORDER BY
CASE status WHEN 'IN_PROGRESS' THEN 0 WHEN 'COMPLETED' THEN 1 ELSE 2 END,
progress_pct ASC

Nombre de lignes et d’octets chargés lors de la mise à jour actuelle

La requête suivante combine le pourcentage de progression, les lignes mises à jour et les octets upload pour chaque table dans l'exécution instantanée actuelle.

SQL
WITH latest_update AS (
SELECT origin.update_id AS update_id
FROM event_log('<pipeline-id>')
WHERE event_type = 'create_update'
ORDER BY timestamp DESC LIMIT 1
),
progress AS (
SELECT
origin.flow_name AS flow_name,
details:operation_progress:progress_percent::double AS progress_pct,
details:operation_progress:status::string AS status,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.update_id = (SELECT update_id FROM latest_update)
),
volume AS (
SELECT
origin.flow_name AS flow_name,
SUM(details:flow_progress:metrics:num_upserted_rows::bigint) AS rows_loaded,
SUM(details:flow_progress:metrics:num_output_bytes::bigint) AS bytes_loaded
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.flow_name LIKE '%_snapshot_flow'
AND origin.update_id = (SELECT update_id FROM latest_update)
GROUP BY origin.flow_name
)
SELECT
p.flow_name,
p.status,
ROUND(p.progress_pct, 2) AS progress_pct,
v.rows_loaded,
ROUND(v.bytes_loaded / 1024.0 / 1024.0, 2) AS mb_loaded,
ROUND(v.bytes_loaded / 1024.0 / 1024.0 / 1024.0, 3) AS gb_loaded
FROM progress p
LEFT JOIN volume v ON p.flow_name = v.flow_name
WHERE p.rn = 1
ORDER BY p.progress_pct ASC

Détection de snapshots bloqués

La query suivante renvoie les tables d'instantanés dont le progress_percent n'a pas changé au cours des 30 dernières minutes. Utilisez-le pour identifier les instantanés bloqués mais toujours actifs.

SQL
WITH latest_update AS (
SELECT origin.update_id AS update_id
FROM event_log('<pipeline-id>')
WHERE event_type = 'create_update'
ORDER BY timestamp DESC LIMIT 1
),
recent AS (
SELECT
origin.flow_name AS flow_name,
details:operation_progress:progress_percent::double AS progress_pct,
details:operation_progress:status::string AS status,
timestamp
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.update_id = (SELECT update_id FROM latest_update)
AND timestamp >= current_timestamp() - INTERVAL 30 MINUTES
)
SELECT
flow_name,
ROUND(MIN(progress_pct), 2) AS min_pct_30min,
ROUND(MAX(progress_pct), 2) AS max_pct_30min,
ROUND(MAX(progress_pct) - MIN(progress_pct), 2) AS pct_change_30min,
COUNT(*) AS events_in_window,
MAX(timestamp) AS last_event_ts
FROM recent
WHERE status = 'IN_PROGRESS'
GROUP BY flow_name
HAVING MAX(progress_pct) - MIN(progress_pct) = 0
AND MAX(progress_pct) < 100
ORDER BY max_pct_30min ASC

ETA de l'instantané par table

La query suivante renvoie le pourcentage de progression actuel et le temps estimé jusqu'à la fin pour chaque table de l'exécution de l'instantané actif. Seules les tables ayant un statut IN_PROGRESS sont renvoyées. Pour inclure les tables terminées et en file d'attente, supprimez le filtre status = 'IN_PROGRESS'.

SQL
WITH latest_update AS (
SELECT origin.update_id AS update_id
FROM event_log('<pipeline-id>')
WHERE event_type = 'create_update'
ORDER BY timestamp DESC LIMIT 1
),
latest_per_flow AS (
SELECT
origin.flow_name AS flow_name,
details:operation_progress:status::string AS status,
details:operation_progress:progress_percent::double AS progress_pct,
details:operation_progress:estimated_completion_ms::bigint AS eta_ms,
timestamp,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log('<pipeline-id>')
WHERE event_type = 'operation_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.update_id = (SELECT update_id FROM latest_update)
)
SELECT
flow_name,
status,
ROUND(progress_pct, 2) AS progress_pct,
eta_ms,
ROUND(eta_ms / 1000.0 / 60.0, 1) AS eta_minutes,
ROUND(eta_ms / 1000.0 / 3600.0, 2) AS eta_hours,
timestamp AS last_update
FROM latest_per_flow
WHERE rn = 1
AND status = 'IN_PROGRESS'
ORDER BY eta_ms DESC NULLS LAST

Querys de latence CDC

Actualité de la CDC par table

Pour chaque flux CDC, la passerelle rapporte plusieurs champs d'observabilité dans streaming_metrics. Ils vous permettent de déterminer si une table est à jour, quelles tables présentent un retard et si le retard se situe côté source ou côté passerelle.

Ce sont des latences de passerelle d'ingestion. Ils mesurent le chemin d'accès de la base de données source à travers la passerelle jusqu'au volume Unity Catalog. Ils n'incluent pas la latence de l'application en aval du volume Unity Catalog vers la table de destination. Cette étape est observée séparément dans le Logs d'événements de l'applier.

Champ

Description

event_time.max

Timestamp ISO 8601 de la modification la plus récente que la passerelle a lue dans la base de données source pour cette table. Si cette valeur cesse de changer entre les événements, la source a cessé de produire des modifications ou la passerelle ne peut plus lire à partir de la source.

discovery_latency_ms

Temps en millisecondes entre le changement de source à event_time.max et le moment où la passerelle a émis cet événement. Couvre le chemin complet de la base de données source au volume Unity Catalog. Une valeur plus faible signifie que les données dans le volume sont plus actuelles.

batch_processing_time_ms

Temps en millisecondes que la passerelle a passé à lire et à upload le dernier batch de modifications. N'inclut pas les retards côté base de données source. Pour estimer le décalage de la base de données source, soustrayez cette valeur de discovery_latency_ms.

Champ

Description

event_time.max

Timestamp ISO 8601 de la modification la plus récente que la passerelle a lue dans la base de données source pour cette table. Si cette valeur cesse de changer entre les événements, la source a cessé de produire des modifications ou la passerelle ne peut plus lire à partir de la source.

discovery_latency_ms

Temps en millisecondes entre le changement de source à event_time.max et le moment où la passerelle a émis cet événement. Couvre le chemin complet de la base de données source au volume Unity Catalog. Une valeur plus faible signifie que les données dans le volume sont plus actuelles.

batch_processing_time_ms

Temps en millisecondes que la passerelle a passé à lire et à upload le dernier batch de modifications. N'inclut pas les retards côté base de données source. Pour estimer le décalage de la base de données source, soustrayez cette valeur de discovery_latency_ms.

La relation entre discovery_latency_ms et batch_processing_time_ms vous indique où, sur le chemin de la passerelle (de la base de données source au volume Unity Catalog), la latence est concentrée :

Modèle

Ce que cela signifie

Les deux petits

Le CDC est en cours d'exécution et à jour.

discovery_latency_ms élevée, batch_processing_time_ms faible

Décalage côté source. La passerelle lit rapidement. Les commit sont arrivés en retard de la source (délai de réplication, arriéré des Logs source, transactions source de longue durée).

Les deux sont élevés

Délai côté passerelle. Le pipeline d'upload est le goulot d'étranglement. Vérifiez les ressources de compute de la passerelle et le chemin réseau vers le volume Unity Catalog.

discovery_latency_ms portant sur une seule table

Problème de source par table (DDL en cours, contention de verrou, emplacement de réplication bloqué, changement de schéma).

discovery_latency_ms en progression sur chaque table

Problème à l’échelle de la passerelle (Ressources, connectivité du volume, accumulation du journal CDC source affectant toutes les captures).

event_time.max ne progressant pas sur plusieurs émissions

La source est inactive, ou la passerelle a perdu la connectivité à la base de données source. Vérifiez l'activation du CDC source et les logs de la passerelle.

Modèle

Ce que cela signifie

Les deux petits

Le CDC est en cours d'exécution et à jour.

discovery_latency_ms élevée, batch_processing_time_ms faible

Décalage côté source. La passerelle lit rapidement. Les commit sont arrivés en retard de la source (délai de réplication, arriéré des Logs source, transactions source de longue durée).

Les deux sont élevés

Délai côté passerelle. Le pipeline d'upload est le goulot d'étranglement. Vérifiez les ressources de compute de la passerelle et le chemin réseau vers le volume Unity Catalog.

discovery_latency_ms portant sur une seule table

Problème de source par table (DDL en cours, contention de verrou, emplacement de réplication bloqué, changement de schéma).

discovery_latency_ms en progression sur chaque table

Problème à l’échelle de la passerelle (Ressources, connectivité du volume, accumulation du journal CDC source affectant toutes les captures).

event_time.max ne progressant pas sur plusieurs émissions

La source est inactive, ou la passerelle a perdu la connectivité à la base de données source. Vérifiez l'activation du CDC source et les logs de la passerelle.

La query suivante renvoie les événements de fraîcheur CDC des 30 dernières minutes pour chaque table CDC. Avec l'intervalle par default de cinq minutes, cela produit environ six lignes par table. Les résultats sont triés par discovery_latency_ms afin que les tables les plus en retard apparaissent en premier.

SQL
SELECT
origin.flow_name AS flow_name,
details:flow_progress:streaming_metrics:event_time:max::string AS latest_source_commit_seen,
details:flow_progress:streaming_metrics:discovery_latency_ms::bigint AS discovery_latency_ms,
details:flow_progress:streaming_metrics:batch_processing_time_ms::bigint AS batch_processing_time_ms,
timestamp AS emission_ts
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.flow_name LIKE '%_cdc_flow'
AND details:flow_progress:streaming_metrics IS NOT NULL
AND timestamp >= current_timestamp() - INTERVAL 30 MINUTES
ORDER BY discovery_latency_ms DESC NULLS LAST, flow_name, emission_ts DESC

Série temporelle de latence CDC

La query suivante renvoie la moyenne horaire de discovery_latency_ms et batch_processing_time_ms pour chaque flux CDC au cours des dernières 24 heures. Utilisez-le pour identifier les moments où l'actualité s'est dégradée. Pour les pipelines avec de nombreux flux CDC, filtrez par nom de flux pour limiter l'ensemble des résultats.

SQL
SELECT
DATE_TRUNC('HOUR', timestamp) AS ts_hour,
origin.flow_name AS flow_name,
ROUND(AVG(details:flow_progress:streaming_metrics:discovery_latency_ms::bigint) / 1000.0, 2) AS avg_discovery_latency_sec,
ROUND(AVG(details:flow_progress:streaming_metrics:batch_processing_time_ms::bigint) / 1000.0, 2) AS avg_batch_processing_sec
FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
AND level = 'METRICS'
AND origin.flow_name IS NOT NULL
AND origin.pipeline_type = 'INGESTION_GATEWAY'
AND origin.flow_name LIKE '%_cdc_flow'
AND details:flow_progress:streaming_metrics IS NOT NULL
AND timestamp >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY DATE_TRUNC('HOUR', timestamp), origin.flow_name
ORDER BY origin.flow_name, ts_hour

Dépannage

Aucun événement de progression n'apparaît

Si vous ne voyez pas les événements de progression dans les Logs d'événements :

  1. Vérifiez que pipelines.gateway.progressEventsEnabled est défini sur "true".
  2. Attendez au moins un intervalle complet après le start du pipeline. default is cinq minutes.
  3. Vérifiez que le pipeline est en cours d'exécution et en ingestion.
  4. Incluez le filtre level = 'METRICS' pour voir uniquement les événements de progression.

Les événements apparaissent trop fréquemment ou trop rarement

Si les événements n'apparaissent pas à la fréquence attendue :

Vérifiez le paramètre pipelines.gateway.progressEventEmitFrequencySeconds et ajustez-le si nécessaire :

  • La valeur default est de cinq minutes (300 secondes).
  • Plage de valeurs valide : de 30 à 3600 secondes. Ajuster au besoin.

Les métriques affichent zéro après le redémarrage du pipeline

Si les métriques se Reset à zéro après le redémarrage d'un pipeline :

Les métriques sont uniquement en mémoire et sont reset lors d'un redémarrage, d'un refresh ou d'une reprise. Ceci est intentionnel pour des raisons de simplicité de mise en œuvre. Le start pipeline accumulera immédiatement de nouvelles métriques.

Métriques manquantes pour certaines tables

Si certaines tables n'affichent pas d'événements de progression :

  1. Assurez-vous que la table n'est pas filtrée dans la configuration du pipeline.
  2. Pour la phase CDC, assurez-vous que la table source a la CDC ou le suivi des modifications activé.
  3. Confirmez que la table est incluse dans la configuration de la passerelle.
  4. Veuillez noter que progress_percent n'est émis que dans les événements operation_progress pour les flux d'instantanés. Les flux CDC n'émettent pas d'événements operation_progress car la CDC n'a pas de concept d'achèvement.

Champs estimated_completion_ms ou streaming_metrics manquants

Si event_log lignes existent mais que l'objet estimated_completion_ms ou streaming_metrics est manquant, consultez Disponibilité des métriques sur les différentes versions de passerelle pour la query de diagnostic et les exigences minimales de la passerelle.

Ressources supplémentaires