REPLACE USING akışlarıyla kısmi anlık görüntü değiştirme

Important

Bu özellik Beta sürümündedir.

REPLACE USING akışı, hedef tabloyu akış kaynağı ile senkronize tutar: belirtilen anahtar sütunlara uyan tüm satırları değiştirir ve diğer tüm verileri değiştirmeden bırakır.

Bir SEQUENCE BY sütunu, güncellemeleri sıraya koyar; böylece güncellemeler yanlış sırada gelse bile sonuç doğru olur. Her anahtar için en yüksek sıra numarası kazanır ve daha düşük sıra numarasına sahip bir satır, hedefte zaten bulunan daha yüksek sıra numaralı bir satırın üzerine asla yazmaz. Aynı anahtara ve aynı sıraya sahip satırlar, değiştirilmek yerine eklenir.

REPLACE USING nasıl çalışır

İki bölge için tıklama ve dönüşüm olaylarını içeren bir olaylar tablosunu düşünün; bunlar şu şekilde seqsıralanır:

region_id device_type event_type seq
1 iOS click 1
1 Android dönüşüm 1
2 iOS click 1
2 Masaüstü click 1

Bir REPLACE USING (region_id) SEQUENCE BY seq akış, 1 ve 3. bölgeler için bu güncellemeleri alır. Bölge 2'de herhangi bir güncelleme yok:

region_id device_type event_type seq
1 iOS click 2
1 Android dönüşüm 2
1 Masaüstü click 2
3 iOS click 1
3 Masaüstü click 2

Hedef şu olur:

region_id device_type event_type seq Outcome
1 iOS click 2 Değiştirildi, çünkü seq 2, seq 1'den büyüktür
1 Android dönüşüm 2 Değiştirildi, çünkü seq 2, seq 1'den büyüktür
1 Masaüstü click 2 Değiştirildi, çünkü 2 sıra numarası 1 sıra numarasından büyüktür
2 iOS click 1 Dokunulmadı, çünkü bu güncellemede anahtar yok
2 Masaüstü click 1 Dokunulmadı, çünkü bu güncellemede anahtar yok
3 Masaüstü click 2 Eklendi. Bölge 3 için seq 1 satırı eklenmez, çünkü bir anahtar için yalnızca en yüksek sıra uygulanır.

Requirements

REPLACE USING akışlarının gereksinimleri şunlardır:

  • REPLACE USING akışları, Databricks Runtime 18.2 ve üzeride, klasik veya sunucusuz hesaplamada çalıştırılır. Databricks, Unity Catalog'u öneriyor.
  • Kaynak bir akış kaynağı olmalıdır. REPLACE USING akış olmayan bir kaynağı reddeder.
  • En az bir anahtar sütun ve tam olarak bir SEQUENCE BY sütun belirtmeniz gerekir.

REPLACE USING akışları ne zaman kullanılır

Lakeflow boru hatları, mevcut satırların üzerine yazan üç akış sunar. Kaynağınızın nasıl göründüğüne ve değiştirilecek satırları nasıl tanımladığına göre seçin:

  • Kaynağınız sütuna göre anahtarlanmış bir dizi kısmi anlık görüntüden oluşuyorsa REPLACE USING kullanın. REPLACE USING, yalnızca gelen verilerle eşleşen verilerin üzerine yazar ve diğer tüm verileri olduğu gibi bırakır. Birincil anahtar gerektirmez.
  • Kaynağınız, açık ekleme, güncelleme ve silme işlemleri içeren bir değişiklik verisi yakalama (CDC) akışı olduğunda veya yavaş değişen boyut (SCD) Tip 2 geçmişine ihtiyaç duyduğunuzda AUTO CDC kullanın. OTO CDC ayrıca gerçek bir birincil anahtar gerektirir. Bkz AUTO CDC API'leri: İşlem hatlarıyla değişiklik verilerini yakalamayı basitleştirin.
  • Kaynağınız bir anlık görüntü olduğunda ve örneğin son 7 gün gibi, bir koşulla seçilen hedef tablodaki bir aralığı toplu işlem olarak yeniden hesaplayıp üzerine yazmak istediğinizde REPLACE WHERE kullanın. Birincil anahtar gerektirmez. Bkz. REPLACE WHERE akışlarıyla toplu işleme.

REPLACE USING akışı oluşturun

SQL veya Python'da REPLACE USING akışlarını tanımlayın.

SQL

FLOW REPLACE USING ile birlikte CREATE STREAMING TABLE satır içi yan tümcesini kullanın:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Alternatif olarak, uzun biçimli CREATE FLOW söz dizimini kullanın:

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME SQL'de gereklidir. Sütunları konumlarına göre değil, adlarına göre eşleştirir.

Python

Tabloyu ve akışı @dp.table ile birlikte tanımlayın:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Alternatif olarak, mevcut bir akış tablosunu @dp.replace_flow ile hedefleyin:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using ana sütunların bir listesidir. sequence_by, bir sütun adı veya Column ifadesidir ve replace_using ayarlandığında gereklidir.

Sıralama ve sırası bozuk veriler

SEQUENCE BY sütunu, sonucu güncellemelerin hangi sırayla geldiğinden bağımsız kılar. Bir satır, yalnızca sıra numarası o anahtar için zaten saklanan sıra numarasından büyükse anahtara uygulanır; bu nedenle, geç gelen veya yeniden oynatılan ve mevcut değerden daha eski olan bir satır göz ardı edilir. Güncellemede bulunmayan anahtarlar değiştirilmeden bırakılır.

Bu uygulamaları takip edin ki değişim öngörülebilir şekilde davransın:

Alıştırma Reason
Anahtar versiyonuna göre kesin olarak artan bir dizi kullanın; örneğin zaman damgası, sürüm numarası veya log ofseti gibi. Aynı anahtara ve aynı sıralamaya sahip iki satır da korunur; bu da söz konusu anahtar için yinelenen satırlar oluşmasına neden olur.
Sıfır olmayan bir dizi kullanın. Bir null dizisi tanımlanmamış davranışa yol açabilir.

Expectations

REPLACE USING akışları beklentileri destekler. warn ve fail diğer akışlarda olduğu gibi davranır: warn satır ihlali sürekli yapar, ihlali kaydeder ve fail güncellemeyi durdurur. Bkz. İşlem hattı beklentileriyle veri kalitesini yönetme.

Bir drop beklenti, ihlalde bulunan bir satırı kaynak onu hiç üretmemiş gibi ele alır. Bırakılan satır, hedef tablodaki eşleşen anahtarların yerine geçmez, onları silmez veya değiştirmez:

  • Ayıklama, tekilleştirmeden önce gerçekleşir; bu nedenle akış, anahtarın en son geçerli sürümünü korur.
  • Bir anahtar için gelen tüm satırlar elenirse, anahtara ait mevcut satırlar olduğu gibi bırakılır.
  • Düşürülen bir satır herhangi bir sıra alt sınırı belirlemediğinden, daha sonra gelen geçerli bir güncelleme, sıra numarası düşürülen satırınkinden daha düşük olsa bile yine de uygulanır.

Sınırlamalar

REPLACE USING akışlarının aşağıdaki sınırlamaları vardır:

  • REPLACE USING hedef tablo başına tek bir akışı destekler. REPLACE USING'i aynı hedefte başka bir akış türüyle birleştirmek desteklenmez.
  • Hedef tablo, işlem hattı içinde oluşturulmalıdır.
  • Kaynak bir akış kaynağı olmalıdır.
  • En az bir anahtar sütun ve bir SEQUENCE BY sütun belirtmeniz gerekir. Ana sütunlar tekrarlanamaz ve her anahtar sütunun türü sıralanabilir olmalıdır. Tam sayılar, dizimler ve tarihler gibi atomik tipler anahtar olabilir, oysa MAP ve VARIANT olamaz.
  • Bağımsız akış tabloları için, sözdizimindeki farklara ilişkin bilgi için REPLACE USING akışlarıyla kısmi anlık görüntü değiştirme uygulama bölümüne bakın.

Examples

Aşağıdaki örnekler, her Unity Kataloğu özellikli çalışma alanında bulunan rezervasyon durumu değişikliklerine dair örnek bir tablodan okunabilir samples.wanderbricks.booking_updates. Her rezervasyon, her değişiklik için bir kez görünür; yani booking_id, yeni bir booking_update_id ile tekrar görünür. Wanderbricks veri setine bakınız.

Örnek 1: Her anahtar için en güncel kaydı tutun

Her rezervasyonun sadece mevcut durumunu sakla. Akış, booking_id alanını anahtar olarak kullanır ve booking_update_id alanına göre sıralanır; bu nedenle, bir rezervasyona ilişkin en güncel güncelleme önceki güncellemelerinin yerini alır. Bunun yerine, kaynağınız ekleme, güncelleme ve silme işlemlerini açıkça içeren bir değişiklik akışıysa AUTO CDC kullanın.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Bu örnek, updated_at zaman damgası yerine booking_update_id öğesine göre sıralanır; çünkü aynı rezervasyona yapılan birden fazla güncelleme aynı zaman damgasını paylaşabilir. Sıralama değeri aynı olan satırlar, değiştirilmek yerine eklenir; bu durumda söz konusu rezervasyonlar için birden fazla satır kalır.

Örnek 2: Birden fazla sütunda anahtar

Bir kayıt sütun kombinasyonuyla tanımlandığında, hepsini içinde REPLACE USINGlisteleyin. Burada her rezervasyon (property_id, booking_id) ile tanımlanır, böylece akış her mülk için her rezervasyonun mevcut durumunu korur. Bir anahtar sütunu null olabiliyorsa, REPLACE USING satırı atlamak yerine null’u null’la eşleştirir.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Örnek 3: Geçersiz kayıtları beklentiyle düşürme

Kötü sıraların hedefe girmemesi için bir beklenti ekleyin. Bırakılan bir satır, kaynak hiç üretmemiş gibi kabul edilir: eşleşen anahtarı değiştirmez veya silmez ve akış o anahtar için en son geçerli satıra geri döner. Bu akış, pozitif bir total_amount olmayan güncellemeleri yok sayar.

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")