Pipeline event log
Le log d'événements de pipeline contient toutes les informations liées à un pipeline, y compris les logs d'audit, les contrôles de qualité des données, la progression du pipeline et la data lineage. Vous pouvez utiliser le log d'événements pour suivre, comprendre et surveiller l'état de vos pipelines de données.
Vous pouvez afficher les entrées du log d’événements dans l'interface utilisateur de monitoring du pipeline, l'API REST de Pipelines ou en interrogeant directement le log d’événements. Cette section se concentre sur l'interrogation directe des Logs d'événements.
Vous pouvez également définir des actions personnalisées à exécuter lorsque des événements sont journalisés, par exemple l'envoi d'alertes, avec les hooks d'événement.
Ne supprimez pas le log d'événements ni le catalogue ou le schéma parent où le log d'événements est publié. La suppression du log d'événements pourrait entraîner l'échec de la mise à jour de votre pipeline lors des exécutions futures.
Pour plus de détails sur le schéma du log des événements du pipeline, consultez le Schéma du log des événements du pipeline.
Consulter les Logs des événements
Cette section décrit le comportement default et la syntaxe pour travailler avec les Logs d'événements pour les pipelines configurés avec Unity Catalog et le mode de publication par default.
- Pour le comportement des pipelines Unity Catalog qui utilisent le mode de publication hérité, consultez Utiliser le journal des événements pour les pipelines en mode de publication hérité Unity Catalog.
- Pour le comportement et la syntaxe des pipelines Hive metastore, consultez Utiliser le journal des événements pour les pipelines Hive metastore.
Par défaut, un pipeline écrit les Logs d'événements dans une table Delta masquée du catalogue et du schéma default configurés pour le pipeline. Bien que masquée, la table peut toujours être interrogée par tous les utilisateurs suffisamment privilégiés. By default, seul l'utilisateur d'exécution du pipeline peut interroger la table des Logs d'événements.
Pour query l’event Logs en tant qu’utilisateur d’exécution, utilisez l’ID de pipeline :
SELECT * FROM event_log(<pipelineId>);
Par default, le nom du Logs d’événement masqué est formaté en event_log_{pipeline_id}, où l’ID du pipeline est l’UUID attribué par le système avec les tirets remplacés par des underscores. La table des Logs d’événements apparaît dans system.information_schema.tables, mais n’est pas visible dans Catalog Explorer ou d’autres pages de l’interface utilisateur du Workspace. Vous devez y accéder à l’aide de la fonction event_log().
Vous pouvez publier le journal des événements en modifiant les **paramètres avancés** de votre pipeline. Pour plus de détails, consultez Paramètres du pipeline pour le log des événements. Lorsque vous publiez un log des événements, spécifiez le nom du log des événements et, facultativement, un catalogue et un schéma, comme dans l'exemple suivant :
{
"id": "ec2a0ff4-d2a5-4c8c-bf1d-d9f12f10e749",
"name": "billing_pipeline",
"event_log": {
"catalog": "catalog_name",
"schema": "schema_name",
"name": "event_log_table_name"
}
}
L'emplacement du Log des événements sert également d'emplacement de schéma pour toutes les queries Auto Loader dans le pipeline. Databricks recommande de créer une vue sur la table du log des événements avant de modifier les privilèges, car certains paramètres de compute pourraient permettre aux utilisateurs d'accéder aux métadonnées de schéma si la table du log des événements est partagée directement. La syntaxe d’exemple suivante crée une vue sur une table de log d’événements et est utilisée dans les exemples de query de log d’événements inclus dans cet article. Remplacez <catalog_name>.<schema_name>.<event_log_table_name> par le nom de table entièrement qualifié de votre Log d'événements de pipeline. Si vous avez publié le log d’événements, utilisez le nom spécifié lors de la publication. Sinon, utilisez event_log(<pipelineId>) où pipelineId est l'ID du pipeline que vous souhaitez query.
CREATE VIEW event_log_raw
AS SELECT * FROM <catalog_name>.<schema_name>.<event_log_table_name>;
Dans Unity Catalog, les vues prennent en charge les requêtes de streaming. L’exemple suivant utilise Structured Streaming pour interroger une vue définie au-dessus d’une table de logs d’événements :
df = spark.readStream.table("event_log_raw")
Exemples de query de base
Les exemples suivants montrent comment interroger le log des événements pour obtenir des informations générales sur les pipelines et pour aider à déboguer des scénarios courants.
Suivez les mises à jour du pipeline en interrogeant les mises à jour précédentes.
L'exemple suivant interroge les mises à jour (ou exécutions ) de votre pipeline, en affichant l'ID de mise à jour, l'état, l'heure de start, l'heure de fin et la durée. Cela vous donne un aperçu des exécutions du pipeline.
Nous supposons que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger le log des événements.
with last_status_per_update AS (
SELECT
origin.pipeline_id AS pipeline_id,
origin.pipeline_name AS pipeline_name,
origin.update_id AS pipeline_update_id,
FROM_JSON(details, 'struct<update_progress: struct<state: string>>').update_progress.state AS last_update_state,
timestamp,
ROW_NUMBER() OVER (
PARTITION BY origin.update_id
ORDER BY timestamp DESC
) AS rn
FROM event_log_raw
WHERE event_type = 'update_progress'
QUALIFY rn = 1
),
update_durations AS (
SELECT
origin.pipeline_id AS pipeline_id,
origin.pipeline_name AS pipeline_name,
origin.update_id AS pipeline_update_id,
-- Capture the start of the update
MIN(CASE WHEN event_type = 'create_update' THEN timestamp END) AS start_time,
-- Capture the end of the update based on terminal states or current timestamp (relevant for continuous mode pipelines)
COALESCE(
MAX(CASE
WHEN event_type = 'update_progress'
AND FROM_JSON(details, 'struct<update_progress: struct<state: string>>').update_progress.state IN ('COMPLETED', 'FAILED', 'CANCELED')
THEN timestamp
END),
current_timestamp()
) AS end_time
FROM event_log_raw
WHERE event_type IN ('create_update', 'update_progress')
AND origin.update_id IS NOT NULL
GROUP BY pipeline_id, pipeline_name, pipeline_update_id
HAVING start_time IS NOT NULL
)
SELECT
s.pipeline_id,
s.pipeline_name,
s.pipeline_update_id,
d.start_time,
d.end_time,
CASE
WHEN d.start_time IS NOT NULL AND d.end_time IS NOT NULL THEN
ROUND(TIMESTAMPDIFF(MILLISECOND, d.start_time, d.end_time) / 1000)
ELSE NULL
END AS duration_seconds,
s.last_update_state AS pipeline_update_status
FROM last_status_per_update s
JOIN update_durations d
ON s.pipeline_id = d.pipeline_id
AND s.pipeline_update_id = d.pipeline_update_id
ORDER BY d.start_time DESC;
Déboguer les problèmes de refresh incrémentiel des vues matérialisées
Cet exemple interroge tous les flux de la dernière mise à jour d'un pipeline. Il indique si elles ont été mises à jour de manière incrémentielle ou non, ainsi que d'autres informations de planification pertinentes qui sont utiles pour le debugging et permettent de comprendre pourquoi un refresh incrémentiel n'a pas lieu.
Nous supposons que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger le log des événements.
WITH latest_update AS (
SELECT
origin.pipeline_id,
origin.update_id AS latest_update_id
FROM event_log_raw AS origin
WHERE origin.event_type = 'create_update'
ORDER BY timestamp DESC
-- LIMIT 1 -- remove if you want to get all of the update_ids
),
parsed_planning AS (
SELECT
origin.pipeline_name,
origin.pipeline_id,
origin.flow_name,
lu.latest_update_id,
from_json(
details:planning_information,
'struct<
technique_information: array<struct<
maintenance_type: string,
is_chosen: boolean,
is_applicable: boolean,
cost: double,
incrementalization_issues: array<struct<
issue_type: string,
prevent_incrementalization: boolean,
operator_name: string,
plan_not_incrementalizable_sub_type: string,
expression_name: string,
plan_not_deterministic_sub_type: string
>>
>>
>'
) AS parsed
FROM event_log_raw AS origin
JOIN latest_update lu
ON origin.update_id = lu.latest_update_id
WHERE details:planning_information IS NOT NULL
),
chosen_technique AS (
SELECT
pipeline_name,
pipeline_id,
flow_name,
latest_update_id,
FILTER(parsed.technique_information, t -> t.is_chosen = true)[0] AS chosen_technique,
parsed.technique_information AS planning_information
FROM parsed_planning
)
SELECT
pipeline_name,
pipeline_id,
flow_name,
latest_update_id,
chosen_technique.maintenance_type,
chosen_technique,
planning_information
FROM chosen_technique
ORDER BY latest_update_id DESC;
Query le coût d'une mise à jour de pipeline
Cet exemple montre comment query l'utilisation des DBU pour un pipeline, ainsi que l'utilisateur pour une exécution de pipeline donnée.
SELECT
sku_name,
billing_origin_product,
usage_date,
collect_set(identity_metadata.run_as) as users,
SUM(usage_quantity) AS `DBUs`
FROM
system.billing.usage
WHERE
usage_metadata.dlt_pipeline_id = :pipeline_id
GROUP BY
ALL;
Requêtes avancées
Les exemples suivants montrent comment interroger le Log d'événements pour gérer des scénarios moins courants ou plus avancés.
query metrics pour tous les flux dans un pipeline
Cet exemple montre comment query des informations détaillées sur chaque flux dans un pipeline. Il affiche le nom du flux, la durée de la mise à jour, les métriques de qualité des données, et des informations sur les lignes traitées (lignes de sortie, supprimées, mises à jour et enregistrements ignorés).
Nous supposons que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger le log des événements.
WITH flow_progress_raw AS (
SELECT
origin.pipeline_name AS pipeline_name,
origin.pipeline_id AS pipeline_id,
origin.flow_name AS table_name,
origin.update_id AS update_id,
timestamp,
details:flow_progress.status AS status,
TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS num_output_rows,
TRY_CAST(details:flow_progress.metrics.num_upserted_rows AS BIGINT) AS num_upserted_rows,
TRY_CAST(details:flow_progress.metrics.num_deleted_rows AS BIGINT) AS num_deleted_rows,
TRY_CAST(details:flow_progress.data_quality.dropped_records AS BIGINT) AS num_expectation_dropped_rows,
FROM_JSON(
details:flow_progress.data_quality.expectations,
SCHEMA_OF_JSON("[{'name':'str', 'dataset':'str', 'passed_records':42, 'failed_records':42}]")
) AS expectations_array
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND origin.flow_name IS NOT NULL
AND origin.flow_name != 'pipelines.flowTimeMetrics.missingFlowName'
),
aggregated_flows AS (
SELECT
pipeline_name,
pipeline_id,
update_id,
table_name,
MIN(CASE WHEN status IN ('STARTING', 'RUNNING', 'COMPLETED') THEN timestamp END) AS start_timestamp,
MAX(CASE WHEN status IN ('STARTING', 'RUNNING', 'COMPLETED') THEN timestamp END) AS end_timestamp,
MAX_BY(status, timestamp) FILTER (
WHERE status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXCLUDED', 'SKIPPED', 'STOPPED', 'IDLE')
) AS final_status,
SUM(COALESCE(num_output_rows, 0)) AS total_output_records,
SUM(COALESCE(num_upserted_rows, 0)) AS total_upserted_records,
SUM(COALESCE(num_deleted_rows, 0)) AS total_deleted_records,
MAX(COALESCE(num_expectation_dropped_rows, 0)) AS total_expectation_dropped_records,
MAX(expectations_array) AS total_expectations
FROM flow_progress_raw
GROUP BY pipeline_name, pipeline_id, update_id, table_name
)
SELECT
af.pipeline_name,
af.pipeline_id,
af.update_id,
af.table_name,
af.start_timestamp,
af.end_timestamp,
af.final_status,
CASE
WHEN af.start_timestamp IS NOT NULL AND af.end_timestamp IS NOT NULL THEN
ROUND(TIMESTAMPDIFF(MILLISECOND, af.start_timestamp, af.end_timestamp) / 1000)
ELSE NULL
END AS duration_seconds,
af.total_output_records,
af.total_upserted_records,
af.total_deleted_records,
af.total_expectation_dropped_records,
af.total_expectations
FROM aggregated_flows af
-- Optional: filter to latest update only
WHERE af.update_id = (
SELECT update_id
FROM aggregated_flows
ORDER BY end_timestamp DESC
LIMIT 1
)
ORDER BY af.end_timestamp DESC, af.pipeline_name, af.pipeline_id, af.update_id, af.table_name;
Métriques de qualité des données ou des attentes de query
Si vous définissez des attentes sur des jeux de données dans votre pipeline, les métriques du nombre d'enregistrements qui ont satisfait ou échoué à une attente sont stockées dans l'objet details:flow_progress.data_quality.expectations. La métrique pour le nombre d'enregistrements ignorés est stockée dans l'objet details:flow_progress.data_quality. Les événements contenant des informations sur la qualité des données ont le type d’événement flow_progress.
Les métriques de qualité des données pourraient ne pas être disponibles pour certains jeux de données. Consultez les limitations d'attente.
Les mesures de qualité des données suivantes sont disponibles :
Métriques | Description |
|---|---|
| Le nombre d'enregistrements qui ont été ignorés parce qu'ils n'ont pas satisfait à une ou plusieurs attentes. |
| Le nombre d'enregistrements qui ont satisfait aux critères d'attente. |
| Le nombre d'enregistrements qui n'ont pas satisfait aux critères d'attente. |
L'exemple suivant interroge les métriques de qualité des données pour la dernière mise à jour du pipeline. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Consulter le Log des événements.
WITH latest_update AS (
SELECT
origin.pipeline_id,
origin.update_id AS latest_update_id
FROM event_log_raw AS origin
WHERE origin.event_type = 'create_update'
ORDER BY timestamp DESC
LIMIT 1 -- remove if you want to get all of the update_ids
),
SELECT
row_expectations.dataset as dataset,
row_expectations.name as expectation,
SUM(row_expectations.passed_records) as passing_records,
SUM(row_expectations.failed_records) as failing_records
FROM
(
SELECT
explode(
from_json(
details:flow_progress:data_quality:expectations,
"array<struct<name: string, dataset: string, passed_records: int, failed_records: int>>"
)
) row_expectations
FROM
event_log_raw,
latest_update
WHERE
event_type = 'flow_progress'
AND origin.update_id = latest_update.id
)
GROUP BY
row_expectations.dataset,
row_expectations.name;
Informations sur la traçabilité des query
Les événements contenant des informations sur la lignée ont le type d'événement flow_definition. L'objet details:flow_definition contient le output_dataset et le input_datasets définissant chaque relation dans le Graphe.
Utilisez la query suivante pour extraire les jeux de données d'entrée et de sortie afin d'afficher les informations de lignage. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Consulter le Log des événements.
with latest_update as (
SELECT origin.update_id as id
FROM event_log_raw
WHERE event_type = 'create_update'
ORDER BY timestamp DESC
limit 1 -- remove if you want all of the update_ids
)
SELECT
details:flow_definition.output_dataset as flow_name,
details:flow_definition.input_datasets as input_flow_names,
details:flow_definition.flow_type as flow_type,
details:flow_definition.schema, -- the schema of the flow
details:flow_definition -- overall flow_definition object
FROM event_log_raw inner join latest_update on origin.update_id = latest_update.id
WHERE details:flow_definition IS NOT NULL
ORDER BY timestamp;
Surveiller l'ingestion de fichiers cloud avec Auto Loader
Les pipelines génèrent des événements lorsque Auto Loader traite des fichiers. Pour les événements Auto Loader, le event_type est operation_progress et le details:operation_progress:type est soit AUTO_LOADER_LISTING soit AUTO_LOADER_BACKFILL. L'objet details:operation_progress comprend également les champs status, duration_ms, auto_loader_details:source_path et auto_loader_details:num_files_listed.
L'exemple suivant interroge les événements Auto Loader pour la dernière mise à jour. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Consulter le Log des événements.
with latest_update as (
SELECT origin.update_id as id
FROM event_log_raw
WHERE event_type = 'create_update'
ORDER BY timestamp DESC
limit 1 -- remove if you want all of the update_ids
)
SELECT
timestamp,
details:operation_progress.status,
details:operation_progress.type,
details:operation_progress:auto_loader_details
FROM
event_log_raw,latest_update
WHERE
event_type like 'operation_progress'
AND
origin.update_id = latest_update.id
AND
details:operation_progress.type in ('AUTO_LOADER_LISTING', 'AUTO_LOADER_BACKFILL');
Surveiller le backlog de données pour optimiser la durée du streaming
Chaque pipeline suit la quantité de données présentes dans le backlog dans l'objet details:flow_progress.metrics.backlog_bytes. Les événements contenant des métriques de backlog ont le type d'événement flow_progress. L'exemple suivant interroge les métriques de backlog pour la dernière mise à jour du pipeline. Cela suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger l'event Logs.
with latest_update as (
SELECT origin.update_id as id
FROM event_log_raw
WHERE event_type = 'create_update'
ORDER BY timestamp DESC
limit 1 -- remove if you want all of the update_ids
)
SELECT
timestamp,
Double(details :flow_progress.metrics.backlog_bytes) as backlog
FROM
event_log_raw,
latest_update
WHERE
event_type ='flow_progress'
AND
origin.update_id = latest_update.id;
Les métriques de backlog pourraient ne pas être disponibles en fonction du type de source de données du pipeline et de la version de Databricks Runtime.
Surveillez les événements d'autoscaling pour optimiser le compute classique
Pour les pipelines qui utilisent le calcul classique (autrement dit, n'utilisent pas le calcul serverless), le Log des événements enregistre les redimensionnements de clusters lorsque la mise à l'échelle automatique améliorée est activée dans vos pipelines. Les événements contenant des informations sur le dimensionnement automatique amélioré ont le type d'événement autoscale. La demande de redimensionnement de cluster informations est stockée dans l'objet details:autoscale.
L'exemple suivant query les demandes de redimensionnement de cluster de dimensionnement automatique amélioré pour la dernière mise à jour du pipeline. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Consulter le Log des événements.
with latest_update as (
SELECT origin.update_id as id
FROM event_log_raw
WHERE event_type = 'create_update'
ORDER BY timestamp DESC
limit 1 -- remove if you want all of the update_ids
)
SELECT
timestamp,
Double(
case
when details :autoscale.status = 'RESIZING' then details :autoscale.requested_num_executors
else null
end
) as starting_num_executors,
Double(
case
when details :autoscale.status = 'SUCCEEDED' then details :autoscale.requested_num_executors
else null
end
) as succeeded_num_executors,
Double(
case
when details :autoscale.status = 'PARTIALLY_SUCCEEDED' then details :autoscale.requested_num_executors
else null
end
) as partially_succeeded_num_executors,
Double(
case
when details :autoscale.status = 'FAILED' then details :autoscale.requested_num_executors
else null
end
) as failed_num_executors
FROM
event_log_raw,
latest_update
WHERE
event_type = 'autoscale'
AND
origin.update_id = latest_update.id
Surveiller l’utilisation des ressources de compute classique
cluster_resources Les événements fournissent des métriques sur le nombre d'emplacements de tâches dans le cluster, leur niveau d'utilisation et le nombre de tâches en attente de planification.
Lorsque le dimensionnement automatique amélioré est activé, les événements cluster_resources contiennent également des métriques pour l'algorithme de dimensionnement automatique, y compris latest_requested_num_executors, et optimal_num_executors. Les événements affichent également l'état de l'algorithme sous différents états tels que CLUSTER_AT_DESIRED_SIZE, SCALE_UP_IN_PROGRESS_WAITING_FOR_EXECUTORS, et BLOCKED_FROM_SCALING_DOWN_BY_CONFIGURATION.
Ces informations peuvent être consultées conjointement avec les événements de dimensionnement automatique pour fournir une vue d'ensemble du dimensionnement automatique amélioré.
L'exemple suivant interroge l'historique de la taille de la file d'attente des tâches, l'historique de l'utilisation, l'historique du nombre d'exécuteurs et d'autres métriques et états pour l'autoscaling lors de la dernière mise à jour du pipeline. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Consulter le Log des événements.
with latest_update as (
SELECT origin.update_id as id
FROM event_log_raw
WHERE event_type = 'create_update'
ORDER BY timestamp DESC
limit 1 -- remove if you want all of the update_ids
)
SELECT
timestamp,
Double(details:cluster_resources.avg_num_queued_tasks) as queue_size,
Double(details:cluster_resources.avg_task_slot_utilization) as utilization,
Double(details:cluster_resources.num_executors) as current_executors,
Double(details:cluster_resources.latest_requested_num_executors) as latest_requested_num_executors,
Double(details:cluster_resources.optimal_num_executors) as optimal_num_executors,
details :cluster_resources.state as autoscaling_state
FROM
event_log_raw,
latest_update
WHERE
event_type = 'cluster_resources'
AND
origin.update_id = latest_update.id;
Surveiller les métriques de streaming de pipeline
Vous pouvez afficher des métriques sur la progression du Stream dans un pipeline. Effectuez une query pour les événements stream_progress afin d'obtenir des événements très similaires aux métriques StreamingQueryListener créées par Structured Streaming, avec les exceptions suivantes :
- Les métriques suivantes sont présentes dans
StreamingQueryListener, mais pas dansstream_progress:numInputRows,inputRowsPerSecondetprocessedRowsPerSecond. - Pour les flux Kafka et Kinesis, les champs
startOffset,endOffsetetlatestOffsetpeuvent être trop volumineux et sont tronqués. Pour chacun de ces champs, un champ...Truncatedsupplémentaire,startOffsetTruncated,endOffsetTruncatedetlatestOffsetTruncated, est ajouté avec une valeur booléenne indiquant si les données sont tronquées.
Pour interroger stream_progress événements, vous pouvez utiliser une query comme la suivante :
SELECT
parse_json(get_json_object(details, '$.stream_progress.progress_json')) AS stream_progress_json
FROM event_log_raw
WHERE event_type = 'stream_progress';
Voici un exemple d’événement, en JSON :
{
"id": "abcd1234-ef56-7890-abcd-ef1234abcd56",
"sequence": {
"control_plane_seq_no": 1234567890123456
},
"origin": {
"cloud": "<cloud>",
"region": "<region>",
"org_id": 0123456789012345,
"pipeline_id": "abcdef12-abcd-3456-7890-abcd1234ef56",
"pipeline_type": "WORKSPACE",
"pipeline_name": "<pipeline name>",
"update_id": "1234abcd-ef56-7890-abcd-ef1234abcd56",
"request_id": "1234abcd-ef56-7890-abcd-ef1234abcd56"
},
"timestamp": "2025-06-17T03:18:14.018Z",
"message": "Completed a streaming update of 'flow_name'."
"level": "INFO",
"details": {
"stream_progress": {
"progress": {
"id": "abcdef12-abcd-3456-7890-abcd1234ef56",
"runId": "1234abcd-ef56-7890-abcd-ef1234abcd56",
"name": "silverTransformFromBronze",
"timestamp": "2022-11-01T18:21:29.500Z",
"batchId": 4,
"durationMs": {
"latestOffset": 62,
"triggerExecution": 62
},
"stateOperators": [],
"sources": [
{
"description": "DeltaSource[dbfs:/path/to/table]",
"startOffset": {
"sourceVersion": 1,
"reservoirId": "abcdef12-abcd-3456-7890-abcd1234ef56",
"reservoirVersion": 3216,
"index": 3214,
"isStartingVersion": true
},
"endOffset": {
"sourceVersion": 1,
"reservoirId": "abcdef12-abcd-3456-7890-abcd1234ef56",
"reservoirVersion": 3216,
"index": 3214,
"isStartingVersion": true
},
"latestOffset": null,
"metrics": {
"numBytesOutstanding": "0",
"numFilesOutstanding": "0"
}
}
],
"sink": {
"description": "DeltaSink[dbfs:/path/to/sink]",
"numOutputRows": -1
}
}
}
},
"event_type": "stream_progress",
"maturity_level": "EVOLVING"
}
Cet exemple montre des enregistrements non tronqués dans une source Kafka, avec les champs ...Truncated définis sur false:
{
"description": "KafkaV2[Subscribe[KAFKA_TOPIC_NAME_INPUT_A]]",
"startOffsetTruncated": false,
"startOffset": {
"KAFKA_TOPIC_NAME_INPUT_A": {
"0": 349706380
}
},
"endOffsetTruncated": false,
"endOffset": {
"KAFKA_TOPIC_NAME_INPUT_A": {
"0": 349706672
}
},
"latestOffsetTruncated": false,
"latestOffset": {
"KAFKA_TOPIC_NAME_INPUT_A": {
"0": 349706672
}
},
"numInputRows": 292,
"inputRowsPerSecond": 13.65826278123392,
"processedRowsPerSecond": 14.479817514628582,
"metrics": {
"avgOffsetsBehindLatest": "0.0",
"estimatedTotalBytesBehindLatest": "0.0",
"maxOffsetsBehindLatest": "0",
"minOffsetsBehindLatest": "0"
}
}
Pipelines d'audit
Vous pouvez utiliser les enregistrements de logs d'événements et d'autres logs d'audit Databricks pour obtenir une image complète de la manière dont les données sont mises à jour dans un pipeline.
LakeFlow Pipelines utilisent les identifiants du propriétaire du pipeline pour exécuter les mises à jour. Vous pouvez modifier les informations d’identification utilisées en changeant le propriétaire du pipeline. Le log d'audit enregistre l'utilisateur pour les actions sur le pipeline, y compris la création de pipelines, les modifications de configuration et le déclenchement des mises à jour.
Voir les événements Unity Catalog pour une référence des événements d'audit Unity Catalog.
Interroger les actions des utilisateurs dans le Logs des événements
Vous pouvez utiliser les Logs des événements pour auditer les événements, par exemple, les actions des utilisateurs. Les événements contenant des informations sur les actions de l'utilisateur sont du type d'événement user_action.
L'information sur l'action est stockée dans l'objet user_action dans le champ details. Utilisez la requête suivante pour construire un journal d'audit des événements utilisateur. Ceci suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger l'event log.
SELECT timestamp, details:user_action:action, details:user_action:user_name FROM event_log_raw WHERE event_type = 'user_action'
|
|
|
|---|---|---|
2021-05-20T19:36:03.517+0000 |
|
|
2021-05-20T19:35:59.913+0000 |
|
|
27.05.2021T00:35:51.971+0000 |
|
|
Informations sur le Runtime
Vous pouvez afficher les informations du runtime pour une mise à jour de pipeline, par exemple la version de Databricks Runtime pour la mise à jour. Cet exemple suppose que vous avez créé la vue event_log_raw pour le pipeline qui vous intéresse, comme décrit dans Interroger le journal des Logs.
SELECT origin.update_id, details:runtime_details:runtime_version:dbr_version FROM event_log_raw WHERE event_type = 'runtime_details'
|
|
|---|---|
| 18,0 |