Surveiller la progression de la passerelle d'ingestion avec les logs d'événements
S’applique à : connecteurs SaaS
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_progressLes é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 deflow_progress. Le pourcentage de progression est cumulatif. Il s'accumule de0à100pendant 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 :
- Accédez à votre passerelle dans le Databricks Workspace.
- Cliquez sur l'onglet Event Logs pour afficher les événements dans l'interface utilisateur.
- 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 :
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:
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é
{
"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
{
"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
{
"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 |
|---|---|---|
| Chaîne | Soit |
| Chaîne | Toujours |
| Chaîne | Timestamp ISO 8601 lorsque l'événement a été généré. |
| Chaîne | Toujours |
| Chaîne | Nom de la passerelle. |
| Chaîne | Nom de la table ingérée. |
| Chaîne | Nom du catalogue Unity Catalog. |
| Chaîne | Nom du schéma Unity Catalog. |
| Chaîne | Identifiant de flux qui indique la phase d'ingestion. Format : |
| Chaîne | Type de base de données source (par exemple, |
| Chaîne | Statut de flux actuel, généralement |
| 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. |
| Entier | Nombre de lignes supprimées depuis le dernier événement. Métrique Delta. Reset à zéro après chaque émission. |
| 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. |
| Entier | Temps en millisecondes entre le changement de source à |
| 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 |
| 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. |
| Chaîne | Type d'opération. |
| Chaîne | Statut actuel de l'opération. |
| Entier | Temps total écoulé de l'opération en millisecondes. |
| Double | Pourcentage d'achèvement de l'instantané ( |
| Entier | Temps estimé restant en millisecondes jusqu'à la fin de l'instantané. Diminue au fur et à mesure que l'instantané progresse et atteint |
| Chaîne | Nom entièrement qualifié de la table faisant l'objet d'une prise d'instantané. |
| Chaîne | Toujours |
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_rowsestnullcar l'instantané ne produit pas de suppressions.
Métriques cumulatives (progress_percent) :
- La valeur s'accumule de
0.0à100.0pendant 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_msetbatch_processing_time_mss'appliquent uniquement aux flux CDC. Les valeurs sont rapportées comme0si le résultat était autrement négatif.event_time.maxS'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_mss'applique uniquement aux flux d'instantanés. La valeur diminue à mesure que l'instantané progresse et atteint0à 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
"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
"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 :
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 :
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.0entre 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
COMPLETEDsignale toujoursprogress_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.
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.
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.
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.
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').
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é.
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.
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.
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.
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.
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.
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.
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'.
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 |
|---|---|
| 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. |
| Temps en millisecondes entre le changement de source à |
| 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 |
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. |
| 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. |
| Problème de source par table (DDL en cours, contention de verrou, emplacement de réplication bloqué, changement de schéma). |
| Problème à l’échelle de la passerelle (Ressources, connectivité du volume, accumulation du journal CDC source affectant toutes les captures). |
| 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.
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.
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 :
- Vérifiez que
pipelines.gateway.progressEventsEnabledest défini sur"true". - Attendez au moins un intervalle complet après le start du pipeline. default is cinq minutes.
- Vérifiez que le pipeline est en cours d'exécution et en ingestion.
- 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 :
- Assurez-vous que la table n'est pas filtrée dans la configuration du pipeline.
- Pour la phase CDC, assurez-vous que la table source a la CDC ou le suivi des modifications activé.
- Confirmez que la table est incluse dans la configuration de la passerelle.
- Veuillez noter que
progress_percentn'est émis que dans les événementsoperation_progresspour les flux d'instantanés. Les flux CDC n'émettent pas d'événementsoperation_progresscar 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.