İşlem hattı beklentileriyle veri kalitesini yönetme

ETL işlem hatlarında akan verileri doğrulayan kalite kısıtlamaları uygulamak için beklentileri kullanın. Beklentiler, veri kalitesi ölçümleri hakkında daha fazla içgörü sağlar ve geçersiz kayıtları algılarken güncelleştirmelerde başarısız olmanıza veya kayıtları bırakmanıza olanak sağlar.

Gelişmiş kullanım örnekleri ve önerilen en iyi yöntemler için bkz . Beklenti önerileri ve gelişmiş desenler.

İşlem hattı beklentileri akış grafiği

Beklentiler nelerdir?

Beklentiler, sorgudan geçen her kayda veri kalitesi kontrolleri uygulayan işlem hattı materyalleştirilmiş görünümde, akış tablosu (streaming table) veya görünüm oluşturma ifadelerinde isteğe bağlı koşullardır. Beklentiler, kısıtlamaları belirtmek için standart SQL Boole deyimlerini kullanır. Tek bir veri kümesi için birden çok beklentiyi birleştirebilir ve bir işlem hattındaki tüm veri kümesi bildirimlerinde beklentileri ayarlayabilirsiniz.

Uyarı

Databricks SQL'de oluşturulan bağımsız bir işlem hattı tarafından desteklenen akış tabloları ve somutlaştırılmış görünümler üzerinde de beklentiler tanımlayabilirsiniz. CONSTRAINT expectation_name EXPECT (expectation_expr) ve CREATE STREAMING TABLE içinde CREATE MATERIALIZED VIEW yan tümcesini kullanın.

Aşağıdaki bölümlerde, bir beklentinin üç bileşeni tanıtılarak söz dizimi örnekleri sağlanır.

Beklenti adı

Her beklentinin, beklentiyi takip etmek ve izlemek için tanımlayıcı olarak kullanılan bir adı olmalıdır. Doğrulanan ölçümleri bildiren bir ad seçin. Aşağıdaki örnek, yaşın 0 ile 120 yaş arasında olduğunu onaylama beklentisini valid_customer_age tanımlar:

Önemli

Belirli bir veri kümesi için bir beklenti adı benzersiz olmalıdır. Bir işlem hattındaki birden çok veri kümesinde beklentileri yeniden kullanabilirsiniz. Bkz . Taşınabilir ve yeniden kullanılabilir beklentiler.

Piton

@dp.table
@dp.expect("valid_customer_age", "age BETWEEN 0 AND 120")
def customers():
  return spark.readStream.table("datasets.samples.raw_customers")

SQL

CREATE OR REFRESH STREAMING TABLE customers(
  CONSTRAINT valid_customer_age EXPECT (age BETWEEN 0 AND 120)
) AS SELECT * FROM STREAM(datasets.samples.raw_customers);

Değerlendirme kısıtlaması

Kısıtlama cümlesi, her kayıt için doğru veya yanlış olarak değerlendirilmesi gereken bir SQL koşullu deyimidir. Kısıtlama, doğrulanan şeyin gerçek mantığını içerir. Bir kayıt bu koşulu sağlamadığında, beklenti tetiklenir.

Kısıtlamalar geçerli SQL söz dizimini kullanmalıdır ve aşağıdakileri içeremez:

  • Özel Python işlevleri
  • Dış hizmet çağrıları
  • Diğer tablolara başvuran alt sorgular

Aşağıda, veri kümesi oluşturma deyimlerine eklenebilen kısıtlamaların örnekleri verilmiştir:

Piton

Python'da kısıtlama söz dizimi şöyledir:

@dp.expect(<constraint-name>, <constraint-clause>)

Birden çok kısıtlama belirtilebilir:

@dp.expect(<constraint-name>, <constraint-clause>)
@dp.expect(<constraint2-name>, <constraint2-clause>)

Examples:

# Simple constraint
@dp.expect("non_negative_price", "price >= 0")

# SQL functions
@dp.expect("valid_date", "year(transaction_date) >= 2020")

# CASE statements
@dp.expect("valid_order_status", """
   CASE
     WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
     WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
     ELSE false
   END
""")

# Multiple constraints
@dp.expect("non_negative_price", "price >= 0")
@dp.expect("valid_purchase_date", "date <= current_date()")

# Complex business logic
@dp.expect(
  "valid_subscription_dates",
  """start_date <= end_date
    AND end_date <= current_date()
    AND start_date >= '2020-01-01'"""
)

# Complex boolean logic
@dp.expect("valid_order_state", """
   (status = 'ACTIVE' AND balance > 0)
   OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
""")

SQL

SQL'de bir kısıtlamanın söz dizimi şöyledir:

CONSTRAINT <constraint-name> EXPECT ( <constraint-clause> )

Birden çok kısıtlama virgülle ayrılmalıdır:

CONSTRAINT <constraint-name> EXPECT ( <constraint-clause> ),
CONSTRAINT <constraint2-name> EXPECT ( <constraint2-clause> )

Examples:

-- Simple constraint
CONSTRAINT non_negative_price EXPECT (price >= 0)

-- SQL functions
CONSTRAINT valid_date EXPECT (year(transaction_date) >= 2020)

-- CASE statements
CONSTRAINT valid_order_status EXPECT (
  CASE
    WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
    WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
    ELSE false
  END
)

-- Multiple constraints
CONSTRAINT non_negative_price EXPECT (price >= 0),
CONSTRAINT valid_purchase_date EXPECT (date <= current_date())

-- Complex business logic
CONSTRAINT valid_subscription_dates EXPECT (
  start_date <= end_date
  AND end_date <= current_date()
  AND start_date >= '2020-01-01'
)

-- Complex boolean logic
CONSTRAINT valid_order_state EXPECT (
  (status = 'ACTIVE' AND balance > 0)
  OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
)

Geçersiz kayıt üzerinde işlem

Kayıt doğrulama denetiminde başarısız olduğunda ne olacağını belirlemek için bir eylem belirtmeniz gerekir. Aşağıdaki tabloda kullanılabilir eylemler açıklanmaktadır:

Eylem SQL söz dizimi Python söz dizimi Result
warn (varsayılan) EXPECT dp.expect Hedefe geçersiz kayıtlar yazılır.
bırakmak EXPECT ... ON VIOLATION DROP ROW dp.expect_or_drop Veriler hedefe yazılmadan önce geçersiz kayıtlar atılır. Bırakılan kayıtların sayısı diğer veri kümesi ölçümleriyle birlikte günlüğe kaydedilir.
başarısız EXPECT ... ON VIOLATION FAIL UPDATE dp.expect_or_fail Geçersiz kayıtlar güncelleştirmenin başarılı olmasını engelliyor. Yeniden işlemeden önce el ile müdahale gerekir.

Ayrıca, veri kaybı veya hata yapmaksızın geçersiz kayıtları karantinaya alabilmeniz için gelişmiş mantık uygulayabilirsiniz. Geçersiz kayıtları karantinaya alma sayfasını görün.

Beklenti izleme ölçümleri

İşlem hattı kullanıcı arabiriminden warn veya drop eylemleri için izleme ölçümlerini görebilirsiniz. fail Geçersiz bir kayıt algılandığında güncelleştirmenin başarısız olmasına neden olduğundan ölçümler kaydedilmez.

Uyarı

Databricks SQL’de oluşturulan bağımsız bir işlem hattı tarafından desteklenen akış tabloları ve somutlaştırılmış görünümler için, işlem hattı kullanıcı arayüzündeki Veri kalitesi sekmesi kullanılamaz. Beklenen ölçümleri görüntülemek için olay günlüğünü sorgular. Bkz Sorgu veri kalitesi veya beklentiler metrikleri.

Beklenti ölçümlerini görüntülemek için aşağıdaki adımları tamamlayın:

  1. Azure Databricks çalışma alanınızın kenar çubuğunda İşler ve İşlem Hatları'na tıklayın.
  2. İşlem hattınızın adına tıklayın.
  3. Beklentileri tanımlanmış bir veri kümesine tıklayın.
  4. Sağ kenar çubuğunda Veri kalitesi sekmesini seçin.

Lakeflow işlem hattı olay günlüğünü sorgulayarak veri kalitesi ölçümlerini görüntüleyebilirsiniz. Bkz Sorgu veri kalitesi veya beklentiler metrikleri.

Geçersiz kayıtları tutma

Beklentilerin varsayılan davranışı geçersiz kayıtların korunmasıdır. Beklentiyi expect ihlal eden ancak kısıtlamayı geçen veya başarısız olan kayıtların sayısını toplamak istediğinizde işlecini kullanın. Beklentiyi ihlal eden kayıtlar, geçerli kayıtlarla birlikte hedef veri kümesine eklenir:

Piton

@dp.expect("valid timestamp", "timestamp > '2012-01-01'")

SQL

CONSTRAINT valid_timestamp EXPECT (timestamp > '2012-01-01')

Geçersiz kayıtları bırakma

Geçersiz kayıtların expect_or_drop daha fazla işlenmesini önlemek için işlecini kullanın. Beklentiyi ihlal eden kayıtlar hedef veri kümesinden bırakılır:

Piton

@dp.expect_or_drop("valid_current_page", "current_page_id IS NOT NULL AND current_page_title IS NOT NULL")

SQL

CONSTRAINT valid_current_page EXPECT (current_page_id IS NOT NULL and current_page_title IS NOT NULL) ON VIOLATION DROP ROW

Geçersiz kayıtlar durumunda başarısız ol

Geçersiz kayıtlar kabul edilemez olduğunda, kayıt doğrulama başarısız olduğunda yürütmeyi expect_or_fail hemen durdurmak için işlecini kullanın. İşlem bir tablo güncelleştirmesiyse, sistem atomik olarak işlemi geri alır:

Piton

@dp.expect_or_fail("valid_count", "count > 0")

SQL

CONSTRAINT valid_count EXPECT (count > 0) ON VIOLATION FAIL UPDATE

Önemli

Tetiklenen işlem hattında tek bir akışın başarısız olması diğer paralel akışların başarısız olmasına neden olmaz. Sürekli bir işlem hattında, bir beklenti başarısızlığı akışı ve ona bağlı tüm akışları durdurur; işlem hattı da neden durduğunu açıklayan bir ileti gönderir.

Bir doğrulama başarısız olduğunda iş akışı düzenlemesi üzerinde daha fazla denetim için, doğrulama ve aşağı akış çalışmalarını ayrı işlem hatlarına bölün ve işlem hattı görevleri arasındaki denetim akışıyla koordine edin. Bkz . Doğrulama tabloları ve işlem hattı denetim akışı.

LFP akış hatası açıklama grafiği

Beklenen başarısız güncelleştirmelerle ilgili sorunları giderme

İşlem hattı, bir beklenti ihlali nedeniyle başarısız olduğunda, işlem hattını yeniden çalıştırmadan önce geçersiz verileri doğru şekilde işlemek için işlem hattı kodunu düzeltmeniz gerekir.

İşlem hatlarının başarısız olması için yapılandırılan beklentiler, ihlalleri algılamak ve bildirmek için gereken bilgileri izlemek için dönüşümlerinizin Spark sorgu planını değiştirir. Bu bilgileri hangi giriş kaydının birçok sorgu için ihlale neden olduğunu belirlemek için kullanabilirsiniz. Lakeflow işlem hatları, bu tür ihlalleri bildirmek için ayrılmış bir hata iletisi sağlar. Aşağıda bir beklenti ihlali hata iletisi örneği verilmişti:

[EXPECTATION_VIOLATION.VERBOSITY_ALL] Flow 'sensor-pipeline' failed to meet the expectation. Violated expectations: 'temperature_in_valid_range'. Input data: '{"id":"TEMP_001","temperature":-500,"timestamp_ms":"1710498600"}'. Output record: '{"sensor_id":"TEMP_001","temperature":-500,"change_time":"2024-03-15 10:30:00"}'. Missing input data: false

Birden çok beklentiyi yönetme

Uyarı

Hem SQL hem de Python tek bir veri kümesinde birden çok beklentiyi desteklese de, yalnızca Python birden çok beklentiyi gruplandırmanıza ve toplu eylemler belirtmenize olanak tanır.

Birden çok beklentiye sahip LFP akış grafiği

Birden çok beklentiyi birlikte gruplandırabilir ve , expect_allve expect_all_or_dropişlevlerini expect_all_or_failkullanarak toplu eylemler belirtebilirsiniz.

Bu dekoratörler, bir Python sözlüğünü argüman olarak kabul eder; burada anahtar, beklenti adı ve değer, beklenti kısıtlamasıdır. İşlem hattınızdaki birden çok veri kümesinde aynı beklenti kümesini yeniden kullanabilirsiniz. Aşağıda Python işleçlerinin her birinin expect_all örnekleri gösterilmektedir:

valid_pages = {"valid_count": "count > 0", "valid_current_page": "current_page_id IS NOT NULL AND current_page_title IS NOT NULL"}

@dp.table
@dp.expect_all(valid_pages)
def raw_data():
  # Create a raw dataset

@dp.table
@dp.expect_all_or_drop(valid_pages)
def prepared_data():
  # Create a cleaned and prepared dataset

@dp.table
@dp.expect_all_or_fail(valid_pages)
def customer_facing_data():
  # Create cleaned and prepared to share the dataset

Sınırlama

  • Yalnızca akış tabloları, gerçekleştirilmiş görünümler ve geçici görünümler beklentileri desteklediği için, veri kalitesi ölçümleri yalnızca bu nesne türleri için desteklenir.
  • Veri kalitesi ölçümleri şu durumlarda kullanılamaz:
    • Sorguda hiçbir beklenti tanımlanmamıştır.
    • Akış, beklentileri desteklemeyen bir işleç kullanır.
    • Sinks gibi akış türleri beklentileri karşılamaz.
    • Belirli bir akış yürütmesi için ilişkili akış tablosunda veya malzeme halinde güncelleştirme yoktur.
    • İşlem hattı yapılandırması, pipelines.metrics.flowTimeReporter.enabled gibi ölçümleri yakalamak için gerekli ayarları içermez.
  • Bazı durumlarda bir COMPLETED akış ölçüm içermeyebilir. Bunun yerine, ölçümler her mikro toplu işlemde durumu flow_progressolan bir RUNNING olayda bildirilir.
  • Görünümler yalnızca sorgulandığında hesaplandığından, tanımlı bir görünümde veri kalitesi ölçümleri kullanılamayabilir. Alternatif olarak, birden çok aşağı akış veri kümesinde sorgulanan bir görünümün birden fazla veri kalitesi ölçüm kümesi olabilir.
  • AUTO CDC FROM SNAPSHOT ile beklentiler desteklenmez.