Azure Databricks'te değişiklik veri akışı özelliğini kullanma

Değişiklik veri akışı (CDF), Delta Lake tablosunun veya Apache Iceberg v3 tablosunun sürümleri arasındaki satır düzeyi değişiklikleri izler. Her değişiklik kaydı, satır verisini ve satırın eklenip eklenmediğini, güncellendiğini veya silindiğini gösteren meta verileri içerir.

Değişiklik veri akışını aşağıdaki yaygın veri kullanım örnekleri için kullanabilirsiniz:

  • İşlem hattının son çalıştırılmasından bu yana yalnızca değişen satırları işleyen artımlı ETL işlem hatları.
  • Uyumluluk ve idare gereksinimleri için veri değişikliklerini izleyen denetim izleri.
  • Aşağı akış tablolarına, önbelleklere veya dış sistemlere yapılan değişiklikleri eşitleyen veri çoğaltma iş yükleri.

Azure Databricks iki yaklaşımı destekler:

  • Otomatik değişiklik veri akışı: Satır kökeni meta verilerini kullanarak tablo okumaları sırasındaki değişiklikleri hesaplar. Bu işlem tek tek tablo yapılandırması gerektirmez ve Delta Lake ve Apache Iceberg v3 tablolarında çalışır. Bkz. Otomatik değişiklik veri akışı.
  • Eski değişiklik veri akışı: Değişiklikleri tabloya yazma işlemleri sırasında somutlaştırır. Yalnızca Delta Lake tablolarını destekler. Tek tek tablo yapılandırması gerektirir. Bkz. Delta Lake için eski değişiklik veri akışı.

Otomatik değişiklik verisi akışı

Important

Bu özellik Genel Önizleme aşamasındadır. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Otomatik değişiklik veri beslemesi, Delta Lake için satır izlemeyi ve Apache Iceberg v3 için satır köken bilgisini kullanarak satır düzeyindeki değişiklikleri yazma zamanında değil, sorgu zamanında hesaplar. Eski değişiklik veri akışının aksine, otomatik değişiklik veri akışı tek tek tablo yapılandırması gerektirmez ve Delta Lake tablolarında ve Apache Iceberg v3 tablolarında çalışır.

MERGE INTO ve UPDATE işlemlerinde her yazma sırasında değişiklikler hesaplanmadığı için, otomatik değişiklik veri beslemesi, eski değişiklik veri beslemesine kıyasla yazma performansını artırır ve depolama maliyetlerini azaltır.

Otomatik değişiklik veri akışı, eski değişiklik veri akışıyla aynı table_changes() ve readChangeFeed API'leri kullanır ve toplu sorgular, Yapılandırılmış Akış ve Databricks-Databricks Delta Lake Sharing ile çalışır. Bkz. Toplu sorgulardaki değişiklikleri okuma ve Değişiklik verilerini artımlı olarak işleme.

Gereksinimler

  • Databricks Runtime 18 LTS veya üzeri
  • Unity Kataloğu'nda kayıtlı desteklenen bir tablo biçimi:
    • Satır izleme etkinleştirilmiş Delta Lake biçiminde veya Iceberg v3 biçiminde yönetilen bir tablo.
    • Satır izleme etkinleştirilmiş Delta Lake biçiminde bir dış tablo.

Bkz. Databricks Unity Kataloğu tablo türleri.

Note

Değişiklik veri akışı Apache Iceberg belirtiminin bir parçası değildir. Azure Databricks okuyucular Apache Iceberg v3 tabloları için otomatik değişiklik veri akışını sorgulayabilir, ancak dış Iceberg okuyucuları bunu yapamaz. Iceberg tablosu spesifikasyonuna bakın.

Delta Lake için yalnızca Azure Databricks okuyucular otomatik değişiklik veri akışını sorgulayabilir.

Değişiklik verileri akışını kullanın

Değişiklik veri akışını kullanmak için gereksinimleri karşılayan bir tablo kullandığınızı doğrulayın. Bkz . Gereksinimler.

Değişiklik veri akışını toplu olarak okumak için aşağıdakileri yapın:

Python

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .table("<table_name>")

Scala

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .table("<table_name>")

SQL

SELECT * FROM table_changes('<table_name>', 0)

Değişiklik veri akışı için toplu okumalar hakkında daha fazla bilgi için bkz. Toplu sorgularda değişiklikleri okuma.

Değişiklik verileri akışını akış olarak okumak için aşağıdakileri yapın:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .table("<table_name>")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .table("<table_name>")

Değişiklik veri akışı için akış okumaları hakkında daha fazla bilgi için bkz. Değişiklik verilerini artımlı olarak işleme.

Eski değişiklik veri akışından geçiş

Delta Lake tablosunu eski değişiklik veri akışından otomatik değişiklik veri akışına geçirmek için aşağıdakileri yapın:

  1. Tablonuzun gereksinimleri karşıladığını doğrulayın.
  2. Aşağıdaki komutu çalıştırarak eski değişiklik veri akışını kapatın:
ALTER TABLE <table_name> UNSET TBLPROPERTIES ('delta.enableChangeDataFeed');

Hem eski hem de otomatik değişiklik veri akışlarını birlikte kullanamazsınız.

Veri akışı şemasını değiştirme

Bir tablonun değişiklik veri akışından okuduğunuzda, sorgu en son tablo sürümü için şemayı kullanır. Azure Databricks çoğu şema değişikliği ve evrim işlemlerini destekler, ancak sütun eşlemesi olan tabloların sınırlamaları vardır. Bkz . Sütun eşlemeli tablolar.

Delta Lake tablosunun şemasındaki veri sütunlarına ek olarak, değişiklik veri akışı değişiklik olayının türünü tanımlayan meta veri sütunları içerir:

Sütun adı Türü Değerler
_change_type String şunları içerir: insert, update_preimage, update_postimage, delete.
preimage güncelleştirmeden önceki değerdir, postimage güncelleştirmeden sonraki değerdir.
_commit_version Long İçerir: Değişikliği içeren Delta günlüğü veya tablo sürümü.
_commit_timestamp Zaman damgası İçerir: İşleme oluşturulduğunda ilişkili zaman damgası.

Şemada bu meta veri sütunlarıyla aynı adlara sahip sütunlar varsa, tablodaki değişiklik veri akışını kullanamazsınız. Veri akışını değiştirme özelliğini açmadan önce, bu çakışmayı çözmek için tablonuzdaki sütunları yeniden adlandırın.

Değişiklik verilerini artımlı olarak işle

Databricks, tablolardaki değişiklikleri artımlı olarak işlemek için değişiklik veri akışını Yapılandırılmış Akış ile birlikte kullanmanızı önerir. Tablonuzun değişiklik veri akışı sürümlerini otomatik olarak izlemek için Azure Databricks için Yapılandırılmış Akış'ı kullanmanız gerekir. SCD tür 1 veya tür 2 tabloları ile CDC işlemesi için bkz. AUTO CDC API'leri: İşlem hatlarıyla değişiklik veri yakalamayı basitleştirme.

Akış ilk kez başladığında, değişiklik verisi akışı tablonun en son anlık görüntüsünü INSERT kayıtları olarak döndürür ve ardından sonraki değişiklikleri değişiklik verisi olarak döndürür. Veri akışlarını değiştirme hem değişiklik verilerini hem de yeni veri satırlarını tablo işlem günlüğüne aynı anda kaydeder.

Bir akışı tablonun değişiklik veri akışını okuyacak şekilde yapılandırmak için seçeneğini readChangeFeedtrue aşağıdaki gibi ayarlayın:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .table("myTable")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .table("myTable")

Hız sınırları

Azure Databricks, hız sınırlarını (maxFilesPerTrigger, maxBytesPerTrigger) ve excludeRegex değişiklik verilerini okurken destekler. Delta Lake akış seçeneklerinin tam listesi için bkz. Delta Lake.

İsteğe bağlı olarak, bir başlangıç sürümü belirtebilirsiniz. Bkz. Başlangıç sürümü belirtme. Başlangıç anlık görüntüsü dışındaki sürümler için hız sınırları işlemelerin tamamına atomik olarak uygulanır. Mevcut toplu iş ya commit’in tamamını içerir ya da commit’i bir sonraki toplu işe erteler.

Tablo geçmişini yeniden yürütme

Değişiklik veri akışı, bir tablodaki tüm değişikliklerin kalıcı bir kaydı olarak hizmet vermek üzere tasarlanmamıştır. Yalnızca değişiklik veri akışı etkinleştirildikten sonra gerçekleşen değişiklikleri kaydeder. Geçerli sürümü ve sonraki tüm değişiklikleri yakalamak için yeni bir akışlı okuma başlatabilirsiniz.

Değişiklik veri akışındaki kayıtlar geçicidir ve yalnızca belirli bir bekletme penceresi için erişilebilir. İşlem günlükleri, tablo sürümlerini ve buna karşılık gelen değişiklik veri akışı sürümlerini düzenli aralıklarla kaldırır. Bir sürüm kaldırıldığında, bu sürüm için değişiklik veri akışını artık okuyasınız.

Kalıcı geçmiş için değişiklik verilerini arşivle

Kullanım örneğiniz bir tablodaki tüm değişikliklerin kalıcı geçmişini tutmanızı gerektiriyorsa, değişiklik veri akışındaki kayıtları yeni bir tabloya yazmak için artımlı mantık kullanın.

Aşağıdaki örnek, mevcut verileri denetim amacıyla veya değişikliklerin baştan sona yeniden oynatılması için toplu iş yükü olarak işlemek üzere trigger.AvailableNow öğesinin nasıl kullanılacağını gösterir:

Python
(spark.readStream
  .option("readChangeFeed", "true")
  .table("source_table")
  .writeStream
  .option("checkpointLocation", <checkpoint-path>)
  .trigger(availableNow=True)
  .toTable("target_table")
)
Scala
spark.readStream
  .option("readChangeFeed", "true")
  .table("source_table")
  .writeStream
  .option("checkpointLocation", <checkpoint-path>)
  .trigger(Trigger.AvailableNow)
  .toTable("target_table")

Başlangıç sürümü belirtme

Belirli bir noktadan gelen değişiklikleri okumak için, zaman damgası veya sürüm numarası kullanarak bir başlangıç sürümü belirtin. Toplu okumalar için başlangıç sürümleri gereklidir. İsteğe bağlı olarak, aralığı sınırlamak için bir bitiş sürümü belirtebilirsiniz. Tablo geçmişi hakkında daha fazla bilgi edinmek için bkz. Zaman yolculuğu.

Değişiklik veri akışını kullanan Yapılandırılmış Akış iş yüklerini yapılandırdığınızda, başlangıç sürümünün belirtilmesi işleme performansını etkileyebilir:

  • Yeni veri işleme işlem hatları genellikle, akış ilk kez başlatıldığında tablodaki tüm mevcut kayıtları işlem olarak INSERT kaydeden varsayılan davranıştan yararlanır.
  • Hedef tablonuz zaten belirli bir noktaya kadar uygun değişiklikleri olan tüm kayıtları içeriyorsa, kaynak tablo durumunun olay olarak INSERT işlenmesini önlemek için bir başlangıç sürümü belirtin.

Aşağıdaki örnekte bozuk bir denetim noktasıyla bir akış hatasından nasıl kurtarılabileceği gösterilmektedir. Bu örnekte, aşağıdaki koşulları varsayın:

  1. Değişiklik veri akışı, tablo oluşturma sırasında kaynak tabloda etkinleştirildi.
  2. Hedef aşağı akış tablosu, sürüm 75 dahil olmak üzere tüm değişiklikleri işledi.
  3. Kaynak tablonun sürüm geçmişi 70 ve üzeri sürümlerde kullanılabilir.

Yazma akışını mevcut hedef tabloya tanımlarken yeni bir denetim noktası konumu belirtmeniz gerekir:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .option("startingVersion", 76)
  .table("source_table")
  .writeStream
  .option("checkpointLocation", "<new-checkpoint-path>")
  .toTable("target_table")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .option("startingVersion", 76)
  .table("source_table")
  .writeStream
  .option("checkpointLocation", "<new-checkpoint-path>")
  .toTable("target_table")

Important

Başlangıç sürümü belirtirseniz ve bu sürüm tablo geçmişinde kullanılamıyorsa, akış yeni bir denetim noktasından başlatılamaz. Yönetilen tablolar geçmiş sürümleri otomatik olarak temizlediğinden, belirtilen tüm başlangıç sürümleri sonunda silinir.

Bkz. Tablo geçmişini yeniden yürütme.

Toplu sorgulardaki değişiklikleri okuma

Toplu sorgu söz dizimini kullanarak belirli bir sürümden başlayan tüm değişiklikleri okuyabilir veya belirtilen sürüm aralığındaki değişiklikleri aşağıdaki gibi okuyabilirsiniz:

  • Sürümleri tamsayılar ve zaman damgaları olarak biçiminde yyyy-MM-dd[ HH:mm:ss[.SSS]]dizeler olarak belirtin.
  • Başlangıç ve bitiş sürümleri dahildir. Başlangıç sürümünden en son sürüme kadar okumak için yalnızca başlangıç sürümünü belirtin.
  • Değişiklik veri akışı etkinleştirilmeden önce bir sürüm belirtirseniz bir hata oluşur.

Toplu okumaları başlangıç ve bitiş sürümü seçenekleriyle kullanmak için aşağıdakileri yapın:

SQL

0 sürümünden 10 sürümüne okumak için aşağıdakileri yapın:

SELECT * FROM table_changes('tableName', 0, 10)

İki zaman damgası sürümü arasında okuma yapmak için aşağıdakileri yapın:

--
SELECT * FROM table_changes('tableName', '2021-04-21 05:45:46', '2021-05-21 12:00:00')

Başlangıç sürümünden en son sürüme kadar okumak için aşağıdakileri yapın:

SELECT * FROM table_changes('tableName', 0)

Adında özel karakterler bulunan bir tablonun değişikliklerini okumak için aşağıdakileri yapın:

SELECT * FROM table_changes('`schema`.`dotted.tableName`', '2021-04-21 06:45:46', '2021-05-21 12:00:00')

Bakınız table_changes tablo değerli fonksiyon.

Python

0 sürümünden 10 sürümüne okumak için aşağıdakileri yapın:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .option("endingVersion", 10) \
  .table("myDeltaTable")

İki zaman damgası arasında okuma yapmak için aşağıdakileri yapın:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingTimestamp", '2021-04-21 05:45:46') \
  .option("endingTimestamp", '2021-05-21 12:00:00') \
  .table("myDeltaTable")

Başlangıç sürümünden en son sürüme kadar okumak için aşağıdakileri yapın:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .table("myDeltaTable")

Scala

0 sürümünden 10 sürümüne okumak için aşağıdakileri yapın:

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .option("endingVersion", 10)
  .table("myDeltaTable")

İki zaman damgası arasında okuma yapmak için aşağıdakileri yapın:

spark.read
  .option("readChangeFeed", "true")
  .option("startingTimestamp", "2021-04-21 05:45:46")
  .option("endingTimestamp", "2021-05-21 12:00:00")
  .table("myDeltaTable")

Başlangıç sürümünden en son sürüme kadar okumak için aşağıdakileri yapın:

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .table("myDeltaTable")

Aralık dışı sürümleri işleme

Varsayılan olarak, son işlemeyi aşan bir sürüm veya zaman damgası belirtirseniz, sorgu hatasını timestampGreaterThanLatestCommitdöndürür.

Databricks Runtime 11.3 LTS ve üzeri sürümlerde, aralık dışı sürümler için toleransı aşağıdaki gibi etkinleştirebilirsiniz:

SET spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = true;

Bu yapılandırma etkinleştirildiğinde sorgu aşağıdaki gibi farklı sonuçlar döndürür:

  • Son işlemenin ötesindeki bir başlangıç sürümü veya zaman damgası boş bir sonuç döndürür.
  • Son işlemenin ötesindeki bir bitiş sürümü veya zaman damgası, başlangıçtan son işlemeye kadar tüm değişiklikleri döndürür.

Delta Lake için eski değişiklik veri akışı

Eski değişiklik veri akışı, tek tek Delta Lake tabloları için el ile yapılandırma gerektirir. Değişiklik veri akışı Apache Iceberg belirtiminde olmadığından Apache Iceberg tabloları desteklenmez. Databricks, otomatik değişiklik veri akışına geçmenizi önerir. Bkz. Eski değişiklik veri akışından geçiş.

Eski değişiklik verisi akışı açıldığında, çalışma zamanı tabloya yazılan tüm veriler için değişiklik olaylarını kaydeder. Bu, belirtilen satırın eklendiğini, silindiğini veya güncelleştirildiğini gösteren meta verilerle birlikte satır verilerini içerir.

Eski değişiklik veri akışı, otomatik değişiklik veri akışıyla aynı readChangeFeed ve table_changes() okuma API'lerini kullanır. Bkz. Değişiklik verilerini artımlı olarak işleme ve Toplu sorgulardaki değişiklikleri okuma.

Eski değişiklik veri akışını açma

Tek tek tablolarda eski değişiklik veri akışını açıkça açmanız gerekir. Aşağıdaki yöntemlerden birini kullanın:

Yeni tablo

komutunda delta.enableChangeDataFeed = true table özelliğini CREATE TABLE ayarlayın.

CREATE TABLE student (id INT, name STRING, age INT)
  TBLPROPERTIES (delta.enableChangeDataFeed = true)

Note

Eski değişiklik veri akışını herhangi bir zaman aralığı için kapatır ve sonra yeniden açarsanız, aralık sorgulanamaz. Aralık boyunca değişiklikleri sorgulamak için otomatik değişiklik veri akışını kullanın. Bkz. Otomatik değişiklik veri akışı.

Varolan tablo

komutunda delta.enableChangeDataFeed = true table özelliğini ALTER TABLE ayarlayın.

ALTER TABLE myDeltaTable
  SET TBLPROPERTIES (delta.enableChangeDataFeed = true)

Depolama ile ilgili dikkat edilmesi gerekenler

Yönetilen tablolar veri değişikliklerini verimli bir şekilde kaydeder ve depolama düzenini iyileştirmek için diğer özellikleri kullanabilir.

Eski değişiklik veri akışında aşağıdaki depolama davranışını göz önünde bulundurmanız gerekir:

  • Değişiklikler ayrı dosyalara kaydedilebileceğinden depolama maliyetlerinde küçük bir artış görebilirsiniz.
  • Yalnızca ekleme veya tam bölüm silme işlemleri gibi bazı işlemler değişiklik veri dosyaları oluşturmaz. Azure Databricks değişiklik veri akışını doğrudan işlem günlüğünden hesaplar.
  • Veri dosyalarını değiştirme, tablonun bekletme ilkesini kullanır. VACUUM komutu değişiklik verisi dosyalarını siler ve işlem günlüğündeki değişiklikler denetim noktası saklama ilkesine tabidir.

Databricks, değişiklik verisi dosyalarını doğrudan sorgulayarak değişiklik veri akışını yeniden yapılandırmaya çalışmamanızı önerir. Delta Lake ve Apache Iceberg API'lerini her zaman kullanın.

Sınırlamalar

Değişiklik veri akışları için aşağıdaki sınırlamaları göz önünde bulundurun:

Sütun eşlemeli tablolar

Delta Lake tablosunda sütun eşleme etkinleştirildiğinde, veri dosyalarını yeniden yazmadan sütunları bırakabilir veya yeniden adlandırabilirsiniz. Delta Lake sütun eşlemesi ile sütunları yeniden adlandırma ve bırakma konusuna bakın.

Ancak, eklemeli olmayan şema değişikliklerinden sonra değişiklik veri akışlarının sınırlamaları vardır. Eklemeli olmayan şema değişiklikleri aşağıdaki işlemleri içerir:

  • Sütunları yeniden adlandırın veya bırakın.
  • Sütun veri türlerini değiştirme.
  • sütun null atanabilirliğini değiştirme, örneğin ile.ALTER COLUMN ... SET NOT NULL Bkz NOT NULL . kısıtlama.

Eklemesiz şema değişikliğinin gerçekleştiği bir işlem veya aralık için değişiklik veri akışlarını okuyamazsınız.

Belirtilen toplu okuma aralığından önce veya sonra eksiz şema değişikliklerine izin vermek için sorgular en son tablo sürümü yerine aralığın son sürümünün şemasını kullanır. Sürüm aralığı eklemeli olmayan bir şema değişikliğine yayılsa bile sorgular başarısız olur.

Otomatik değişiklik veri beslemesi

  • Değişiklik veri akışı Apache Iceberg belirtiminde desteklenmediğinden, dış Iceberg istemcileri otomatik değişiklik veri akışını sorgulayamaz. Iceberg tablosu spesifikasyonuna bakın.
  • Birden çok deyim içeren işlemlerde, kaynak tablo işlem sırasında değiştirilmişse otomatik değişiklik verisi akışı desteklenmez.
  • Satır filtreleri veya sütun maskeleri olan tablolarda otomatik değişiklik veri akışı desteklenmez. Bkz. Satır filtreleri ve sütun maskeleri.
  • Değişiklik veri akışı sorguları, sütun yeniden adlandırma, bırakma veya veri türü değişikliği gibi eklemesiz şema değişikliğinin gerçekleştiği tablo sürümlerine yayılamaz. Sorguyu şema değişikliğinden önce ve sonra aralıklara bölün.