Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
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 BYsü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 BYsü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, oysaMAPveVARIANTolamaz. - 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")