Pipeline olay günlüğü

İşlem hattı olay günlüğü, denetim günlükleri, veri kalitesi denetimleri, işlem hattı ilerleme durumu ve veri kökeni dahil olmak üzere işlem hattıyla ilgili tüm bilgileri içerir. Veri işlem hatlarınızın durumunu izlemek, anlamak ve izlemek için olay günlüğünü kullanabilirsiniz.

İşlem hattı izleme kullanıcı arabiriminde, pipelines REST API'sinde veya olay günlüğünü doğrudan sorgulayarak olay günlüğü girdilerini görüntüleyebilirsiniz. Bu bölüm, olay günlüğünü doğrudan sorgulamaya odaklanır.

Olaylar günlüğe kaydedilirken, örneğin uyarılar göndermek gibi, özel eylemler çalıştıracak olay kancalarıtanımlayabilirsiniz.

Önemli

Olay günlüğünü, olay günlüğünün yayımlandığı üst kataloğu veya şemayı silmeyin. Olay günlüğünü silmek, işlem hattınızın gelecekteki çalıştırmalar sırasında güncelleştirilememesine neden olabilir.

Olay günlüğü şemasının tüm ayrıntıları için bkz. İşlem hattı olay günlüğü şeması.

Olay günlüğünü sorgulama

Uyarı

Bu bölümde Unity Kataloğu ve varsayılan yayımlama modu ile yapılandırılmış işlem hatları için olay günlükleriyle çalışmaya yönelik varsayılan davranış ve söz dizimi açıklanmaktadır.

Varsayılan olarak, işlem hattı olay günlüğünü işlem hattı için yapılandırılan varsayılan katalog ve şemada gizli bir Delta tablosuna yazar. Gizliyken, tablo yine de tüm yeterince ayrıcalıklı kullanıcılar tarafından sorgulanabilir. Varsayılan olarak, olay günlüğü tablosunu yalnızca işlem hattının adına çalıştırılan kullanıcısı sorgulayabilir.

Olay günlüğünü kimliğine bürünülen kullanıcı olarak sorgulamak için işlem hattı kimliğini kullanın:

SELECT * FROM event_log(<pipelineId>);

Varsayılan olarak, gizli olay günlüğünün adı event_log_{pipeline_id} olarak biçimlendirilir; burada pipeline kimliği, tirelerin alt çizgilerle değiştirildiği sistem tarafından atanan UUID'dir. Olay günlüğü tablosu içinde system.information_schema.tables görünür, ancak Katalog Gezgini'nde veya diğer çalışma alanı kullanıcı arabirimi sayfalarında görünmez. event_log() işlevini kullanarak ona erişmelisiniz.

İşlem hattınızın Gelişmiş ayarlarını düzenleyerek olay günlüğünü yayımlayabilirsiniz. Ayrıntılar için bkz. Olay günlüğü için akış ayarları. Bir olay günlüğü yayımladığınızda, olay günlüğünün adını belirtin ve isteğe bağlı olarak aşağıdaki örnekte olduğu gibi bir katalog ve şema belirtin:

{
  "id": "ec2a0ff4-d2a5-4c8c-bf1d-d9f12f10e749",
  "name": "billing_pipeline",
  "event_log": {
    "catalog": "catalog_name",
    "schema": "schema_name",
    "name": "event_log_table_name"
  }
}

Olay günlüğü konumu, işlem hattındaki tüm Otomatik Yükleyici sorguları için şema konumu olarak da hizmet eder. Databricks, ayrıcalıkları değiştirmeden önce olay günlüğü tablosu üzerinde bir görünüm oluşturmanızı önerir. Bazı işlem ayarları, olay günlüğü tablosu doğrudan paylaşılırsa kullanıcıların şema meta verilerine erişmesine izin verebilir. Aşağıdaki örnek söz dizimi, olay günlüğü tablosunda bir görünüm oluşturur ve bu makalede yer alan örnek olay günlüğü sorgularında kullanılır. İşlem hattı olay günlüğünüzün tam tablo adıyla <catalog_name>.<schema_name>.<event_log_table_name> değerini değiştirin. Olay günlüğünü yayımladıysanız yayımlarken belirtilen adı kullanın. Aksi takdirde, pipelineId değerinin sorgulamak istediğiniz işlem hattının kimliği olduğu yeri kullanın event_log(<pipelineId>) .

CREATE VIEW event_log_raw
AS SELECT * FROM <catalog_name>.<schema_name>.<event_log_table_name>;

Unity Kataloğu'nda görünümler akış sorgularını destekler. Aşağıdaki örnek, bir olay günlüğü tablosunun üzerinde tanımlanan bir görünümü sorgulamak için Yapılandırılmış Akış'ı kullanır:

df = spark.readStream.table("event_log_raw")

Temel sorgu örnekleri

Aşağıdaki örneklerde işlem hatları hakkında genel bilgi almak ve yaygın senaryolarda hata ayıklamaya yardımcı olmak için olay günlüğünü sorgulama işlemi gösterilmektedir.

Önceki güncelleştirmeleri sorgulayarak işlem hattı güncelleştirmelerini izleme

Aşağıdaki örnek işlem hattınızın güncelleştirmelerini (veya çalıştırmalarını) sorgular ve güncelleştirme kimliğini, durumunu, başlangıç saatini, tamamlanma süresini ve süreyi gösterir. Bu, işlem hattı çalıştırmaları hakkında genel bir bakış sağlar.

event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

Gerçekleştirilmiş görünümün artımlı yenileme sorunlarını ayıklama

Bu örnek, bir işlem hattının en son güncelleştirmesinden gelen tüm akışları sorgular. Artımlı olarak güncelleştirilip güncelleştirilmediğini ve artımlı yenileme sorununu hata ayıklamak için yararlı olan diğer ilgili planlama bilgilerini gösterir.

event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

İşlem hattı güncelleştirmesinin maliyetini sorgulama

Bu örnek, bir işlem hattı için DBU kullanımını ve belirli bir işlem hattı çalıştırması için kullanıcıyı sorgulamayı gösterir.

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;

Gelişmiş sorgular

Aşağıdaki örneklerde daha az yaygın veya daha gelişmiş senaryoları işlemek için olay günlüğünü sorgulama gösterilmektedir.

İşlem hattındaki tüm akışlar için sorgu ölçümleri

Bu örnekte, bir işlem hattındaki her akışla ilgili ayrıntılı bilgilerin nasıl sorgulandığı gösterilmektedir. Akış adını, güncelleştirme süresini, veri kalitesi ölçümlerini ve işlenen satırlarla ilgili bilgileri (çıkış satırları, silinmiş, yükseltilmiş ve bırakılan kayıtlar) gösterir.

event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

Veri kalitesi veya beklentileri ölçümlerini sorgulama

İşlem hattınızdaki veri kümeleriyle ilgili beklentileri tanımlarsanız, bir beklentiyi geçen ve başarısız olan kayıt sayısına ilişkin ölçümler nesnesinde details:flow_progress.data_quality.expectations depolanır. Düşen kayıtların sayısını belirten ölçüm, details:flow_progress.data_quality nesnesinde saklanır. Veri kalitesi hakkında bilgi içeren olaylar flow_progressolay türüne sahiptir.

Veri kalitesi ölçümleri bazı veri kümelerinde kullanılamayabilir. Beklenti sınırlamalarına bakın.

Aşağıdaki veri kalitesi ölçümleri kullanılabilir:

Ölçü birimi Description
dropped_records Bir veya daha fazla beklentiyi başarısız olduğu için düşürülen kayıtların sayısı.
passed_records Beklenti ölçütlerini geçen kayıtların sayısı.
failed_records Beklenti ölçütlerini karşılamayan kayıtların sayısı.

Aşağıdaki örnek, son işlem hattı güncelleştirmesi için veri kalitesi ölçümlerini sorgular. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

Sorgu kökeni bilgileri

Köken hakkında bilgi içeren olaylar flow_definitionolay türüne sahiptir. details:flow_definition nesnesi output_dataset içerir ve input_datasets grafikteki her ilişkiyi tanımlar.

Köken bilgilerini görmek için giriş ve çıkış veri kümelerini ayıklamak için aşağıdaki sorguyu kullanın. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

Otomatik Yükleyici ile bulut dosyası alımını izleme

İşlem hatları, Otomatik Yükleyici dosyaları işlediğinde olaylar oluşturur. Otomatik Yükleyici olayları için event_typeoperation_progress ve details:operation_progress:type, ya AUTO_LOADER_LISTING ya da AUTO_LOADER_BACKFILL'dür. details:operation_progress nesnesi status, duration_ms, auto_loader_details:source_pathve auto_loader_details:num_files_listed alanlarını da içerir.

Aşağıdaki örnek, en son güncelleştirme için Otomatik Yükleyici olaylarını sorgular. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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');

Akış süresini iyileştirmek için veri birikimini izleme

Her işlem hattı, details:flow_progress.metrics.backlog_bytes nesnesinde bulunan geri dönüşümde ne kadar veri olduğunu takip eder. Geribildirim ölçümlerini içeren olaylar, olay türü flow_progressolanlardır. Aşağıdaki örnek, son pipeline güncellemesi için birikim ölçümlerini sorgular. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

Uyarı

İşlem hattının veri kaynağı türüne ve Databricks Runtime sürümüne bağlı olarak, birikim ölçümleri kullanılamayabilir.

Klasik işlemi iyileştirmek için otomatik ölçeklendirme olaylarını izleme

Klasik işlem kullanan işlem hatları için (başka bir deyişle sunucusuz işlem kullanmayın), işlem hatlarınızda gelişmiş otomatik ölçeklendirme etkinleştirildiğinde olay günlüğü kümeyi yeniden boyutlandırıyor. Gelişmiş otomatik ölçeklendirme hakkında bilgi içeren olayların olay türü autoscale. Küme yeniden boyutlandırma isteği bilgileri details:autoscale nesnesinde depolanır.

Aşağıdaki örnek, son işlem hattı güncelleştirmesi için gelişmiş otomatik ölçeklendirme kümesi yeniden boyutlandırma isteklerini sorgular. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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

Klasik işlem için işlem kaynağı kullanımını izleme

cluster_resources olaylar, kümedeki görev yuvalarının sayısı, bu görev yuvalarının ne kadarının kullanıldığı ve zamanlamayı bekleyen görev sayısıyla ilgili ölçümler sağlar.

İyileştirilmiş otomatik ölçeklendirme etkinleştirildiğinde, cluster_resources olayları latest_requested_num_executorsve optimal_num_executorsdahil olmak üzere otomatik ölçeklendirme algoritması için ölçümler de içerir. Olaylar ayrıca algoritmanın durumunu CLUSTER_AT_DESIRED_SIZE, SCALE_UP_IN_PROGRESS_WAITING_FOR_EXECUTORSve BLOCKED_FROM_SCALING_DOWN_BY_CONFIGURATIONgibi farklı durumlar olarak gösterir. Bu bilgiler, gelişmiş otomatik ölçeklendirmenin genel bir resmini sağlamak için otomatik ölçeklendirme olaylarıyla birlikte görüntülenebilir.

Aşağıdaki örnek, son işlem hattı güncelleştirmesinde otomatik ölçeklendirme için görev kuyruğu boyut geçmişini, kullanım geçmişini, yürütücü sayısı geçmişini ve diğer ölçümleri ve durumu sorgular. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

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;

İşlem hattı akış ölçümlerini izleme

İşlem hattında akış ilerleme durumuyla ilgili ölçümleri görüntüleyebilirsiniz. Aşağıdaki özel durumlar dışında Yapılandırılmış Akış tarafından oluşturulan StreamingQueryListener ölçümlerine çok benzeyen olayları almak için stream_progress olayları sorgulayın:

  • Aşağıdaki ölçümler içinde StreamingQueryListenerbulunur, ancak içinde stream_progressyoktur: numInputRows, inputRowsPerSecondve processedRowsPerSecond.
  • Kafka ve Kinesis akışları için startOffset, endOffset ve latestOffset alanları çok büyük olabilir ve kısaltılmıştır. Bu alanların her biri için, verilerin kesilip kesilmediğini belirtmek amacıyla Boole değeri taşıyan ek bir ...Truncated alanı, startOffsetTruncated, endOffsetTruncated ve latestOffsetTruncated eklenir.

Olayları sorgulamak için stream_progress aşağıdaki gibi bir sorgu kullanabilirsiniz:

SELECT
  parse_json(get_json_object(details, '$.stream_progress.progress_json')) AS stream_progress_json
FROM event_log_raw
WHERE event_type = 'stream_progress';

JSON'da bir olay örneği aşağıda verilmiştir:

{
  "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"
}

Bu örnekte, bir Kafka kaynağında, ...Truncated alanları false olarak ayarlanan kesilmemiş kayıtlar gösterilir.

{
  "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"
  }
}

İşlem hatlarını denetleme

İşlem hattında verilerin nasıl güncelleştirilmekte olduğunu tam olarak görmek için olay günlüğü kayıtlarını ve diğer Azure Databricks denetim günlüklerini kullanabilirsiniz.

Lakeflow işlem hatları, güncelleştirmeleri çalıştırmak için işlem hattı sahibinin kimlik bilgilerini kullanır. İşlem hattı sahibini değiştirerek kullanılan kimlik bilgilerini değiştirebilirsiniz. Denetim günlüğü işlem hattı oluşturma, yapılandırma düzenleme ve güncelleştirmeleri tetikleme dahil olmak üzere işlem hattındaki eylemler için kullanıcıyı kaydeder.

Unity Kataloğu denetim olayları için bir referans olarak Unity Kataloğu olaylarını görün.

Olay günlüğünde kullanıcı eylemlerini sorgulama

Olay günlüğünü kullanarak olayları, örneğin kullanıcı eylemlerini denetleyebilirsiniz. Kullanıcı eylemleri hakkında bilgi içeren olaylar user_actionolay türüne sahiptir.

Eylem hakkındaki bilgiler user_action alanındaki details nesnesinde depolanır. Kullanıcı olaylarının denetim günlüğünü oluşturmak için aşağıdaki sorguyu kullanın. Bu, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuzu varsayar.

SELECT timestamp, details:user_action:action, details:user_action:user_name FROM event_log_raw WHERE event_type = 'user_action'
timestamp action user_name
2021-05-20T19:36:03.517+0000 START user@company.com
2021-05-20T19:35:59.913+0000 CREATE user@company.com
2021-05-27T00:35:51.971+0000 START user@company.com

Çalışma zamanı bilgileri

İşlem hattı güncelleştirmesinin çalışma zamanı bilgilerini, örneğin güncelleştirmenin Databricks Runtime sürümünü görüntüleyebilirsiniz. Bu örnekte, event_log_raw bölümünde açıklandığı gibi ilgilendiğiniz işlem hattı için görünümü oluşturduğunuz varsayılır.

SELECT origin.update_id, details:runtime_details:runtime_version:dbr_version FROM event_log_raw WHERE event_type = 'runtime_details'
update_id dbr_version
1234abcd-ef56-7890-abcd-ef1234abcd56 18,0