Lakebase'e bağlanma

Lakebase'e veya dahili toplu işleme, otomatik denemeler ve çalışma alanı tarafından yönetilen kimlik doğrulama ile harici bir PostgreSQL veritabanına yazmak için Structured Streaming kullanın.

Lakebase havuzu ne zaman kullanılır?

Lakebase’e veya harici bir PostgreSQL veritabanına düşük gecikmeli akış yazma işlemleri için Lakebase alıcısını kullanın. Bu alıcı, toplu işleme, bağlantı yönetimi ve hata işleme için özel foreach işlevleri uygulamanızı gerektirmez.

Yaygın kullanım örnekleri şunlardır:

  • İşletimsel panolar veya müşteriye yönelik özellikler için uygulama veritabanlarını gerçek zamanlı olarak güncelleştirin.
  • Toplu veya filtrelenmiş akış sonuçları gibi sürekli değişen verileri işlemsel bir veritabanına eşitleyin.
  • Gerçek zamanlı modu kullanarak bir Yapılandırılmış Akış sorgusunun çıkışını bir Lakebase tablosuna saniyenin altında gecikme süresiyle yazın.

Lakebase'den Lakehouse içindeki Delta Lake tablolarına veri eşitlemek için, yani ters yönde, bkz. Lakebase Değişiklik Veri Akışı.

Gereksinimler

  • Databricks Runtime 18 LTS ve üzeri.
    • Harici PostgreSQL bağlantıları için Databricks Runtime 19 ve üzeri kullanmanız ve UC Compute önizlemesinde Özel JDBC'ye katılmanız gerekir.
    • Aralıklı veri türleri için Databricks Runtime 19 ve üzeri sürümleri kullanmanız gerekir.
  • Özel veya standart erişim modlarına sahip klasik hesaplama veya dizüstü bilgisayarlar veya işler için sunucusuz hesaplama. Sunucusuz bilgi işlemde Trigger.AvailableNow() kullanın. Bkz. Sunucusuz işlemde akış.
  • Bir Lakebase veritabanı veya harici bir PostgreSQL veritabanına Unity Catalog bağlantısı.

Tanımlayıcı gereksinimleri

Tüm hedefler için Databricks, yalnızca harf, sayı ve alt çizgi içeren şema, tablo, sütun ve birincil anahtar sütun adlarının kullanılmasını önerir. Sink, otomatik olarak bir Lakebase tablosu oluşturduğunda bu gereksinimleri uygular. Bu gereksinimleri karşılamayan tanımlayıcıları kullanmak için, sorguya başlamadan önce hedef tablo oluşturun.

Veritabanına bağlan

Lakebase havuzu aşağıdaki bağlantı yöntemlerini destekler:

Unity Kataloğu'yla kaydedilen Lakebase tabloları

Unity Kataloğu'na kayıtlı Lakebase tabloları için bağlayıcı kimlik bilgilerini otomatik olarak yönetir ve sorguyu çalıştıran kullanıcı veya hizmet sorumlusunun kimliğini kullanır. Tablo yoksa, bağlayıcı tabloyu oluşturur.

Bir Lakebase veritabanını Unity Kataloğu'na kaydetmek için bkz. Unity Kataloğu'nda Lakebase veritabanını kaydetme.

Bir Lakebase tablosuna yazmak için, .toTable() tam nitelikli bir tablo adına sahip yöntemi kullanın: catalog.schema.table

Python

(df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")
)

Scala

df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")

Aşağıdaki yer tutucuları değiştirin:

  • <catalog>.<schema>.<table>: Hedef tablonun tam adı. catalog, Lakebase veritabanını kaydettiğinizde oluşturduğunuz Unity Kataloğu kataloğudur. Bkz. Unity Kataloğu'nda Bir Lakebase veritabanını kaydetme. Tablo yoksa, bağlayıcı bunu oluşturur.
  • <primary-key-columns>: İsteğe bağlı. Hedef tablonun birincil anahtarındaki tüm sütunların virgülle ayrılmış bir listesi, örneğin id veya user_id,event_type. Upsert davranışına bakınız.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Sorgunun denetim noktası verilerini depoladığı bir Unity Catalog birim yolu. Bulut nesnesi depolama URI'sini de kullanabilirsiniz. Konum, yerel diskte değil, yazabileceğiniz depolama alanı olmalı ve her akış sorgusu için benzersiz olmalıdır. Bu, hedef tablodan bağımsızdır. Bkz . Yapılandırılmış Akış denetim noktaları.

Isteğe bağlı yapılandırmalar için, örneğin batchsize ve batchinterval, PostgreSQL sink seçeneklerine bakınız.

Unity Kataloğu'yla kaydedilmeyen Lakebase tabloları

Unity Kataloğu'na kayıtlı olmayan Lakebase tabloları için bağlayıcı kimlik bilgilerini otomatik olarak yönetir ve sorguyu çalıştıran kullanıcı veya hizmet sorumlusunun kimliğini kullanır. Tablo yoksa, bağlayıcı tabloyu oluşturur.

Bir Lakebase tablosuna yazmak için dbtable ve endpoint seçeneklerini kullanın:

Python

(df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
  .option("database", "<database>")  # Optional. Defaults to databricks_postgres.
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()
)

Scala

df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
  .option("database", "<database>")  // Optional. Defaults to databricks_postgres.
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()

Aşağıdaki yer tutucuları değiştirin:

  • <project-id>.<branch-id>.<endpoint-id>: Lakebase uç noktanız. Computes sekmesindeki Kimlik Al menüsünde, biçimi olan projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id> içinde üç değerin tümünü bulabilirsiniz. Bkz . İşlem tanımlayıcıları.
  • <database>: İsteğe bağlı. Hedef PostgreSQL veritabanının adı. Varsayılan değer databricks_postgres’dır. Bkz. Veritabanlarını yönetme.
  • <schema>.<table>: schema.table biçimindeki hedef tablo. Şemayı atlarsanız havuz şemayı public kullanır. Otomatik tablo oluşturma için, harf veya alt çizgi ile başlayan ve yalnızca harf, sayı ve alt çizgi içeren tanımlayıcılar kullanın.
  • <primary-key-columns>: İsteğe bağlı. Hedef tablonun birincil anahtarındaki tüm sütunların virgülle ayrılmış bir listesi, örneğin id veya user_id,event_type. Upsert davranışına bakınız.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Sorgunun denetim noktası verilerini depoladığı bir Unity Catalog birim yolu. Bulut nesnesi depolama URI'sini de kullanabilirsiniz. Konum, yerel diskte değil, yazabileceğiniz depolama alanı olmalı ve her akış sorgusu için benzersiz olmalıdır. Bu, hedef tablodan bağımsızdır. Bkz . Yapılandırılmış Akış denetim noktaları.

Isteğe bağlı yapılandırmalar için, örneğin batchsize ve batchinterval, PostgreSQL sink seçeneklerine bakınız.

Unity Kataloğu kimlik bilgileriyle Harici PostgreSQL

Important

Bu özellik Genel Önizleme aşamasındadır. Çalışma alanı yöneticileri, Önizlemeler sayfasından UC Compute'da Özel JDBC'ye erişimi kontrol edebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Kodunuza kimlik bilgileri kaydetmeden harici bir PostgreSQL veritabanına kimlik doğrulama yapmak için Unity Kataloğu bağlantısı kullanın. Hedef tablo zaten var olmalı.

Bir tür POSTGRESQLbağlantı oluşturun, bkz. Bağlantı oluşturun. Sorguyu çalıştıran kullanıcı veya hizmet prensibi bağlantıda olmalıdır USE CONNECTION .

PostgreSQL tablosuna yazmak için , database, ve dbtable seçenekler kullanındatabricks.connection:

Python

(df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("databricks.connection", "<connection-name>")
  .option("database", "<database>")
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()
)

Scala

df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("databricks.connection", "<connection-name>")
  .option("database", "<database>")
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()

Aşağıdaki yer tutucuları değiştirin:

  • <connection-name>: Unity Kataloğu bağlantısının adı.
  • <database>: Hedef PostgreSQL veritabanının adı.
  • <schema>.<table>: schema.table biçimindeki mevcut hedef tablo. Şemayı atlarsanız havuz şemayı public kullanır.
  • <primary-key-columns>: İsteğe bağlı. Hedef tablonun birincil anahtarındaki tüm sütunların virgülle ayrılmış bir listesi, örneğin id veya user_id,event_type. Upsert davranışına bakınız.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Sorgunun denetim noktası verilerini depoladığı bir Unity Catalog birim yolu. Bulut nesnesi depolama URI'sini de kullanabilirsiniz. Konum, yerel diskte değil, yazabileceğiniz depolama alanı olmalı ve her akış sorgusu için benzersiz olmalıdır. Bu, hedef tablodan bağımsızdır. Bkz . Yapılandırılmış Akış denetim noktaları.

PostgreSQL bağlantıları her zaman TLS kullanır. Sertifika doğrulaması, bağlantıyı oluştururken seçtiğiniz Unity Kataloğu bağlantısındaki ayarları takip eder:

  • Güven sunucusu sertifikası: Seçildiğinde, bağlantı sslmode=require, kullanır ve bu da sunucu sertifikasını doğrulamadan bağlantıyı şifreler.
  • Kullanıcı tarafından sağlanan sunucu sertifikası: Trust sunucu sertifikası seçilmediğinde kullanmak sslmode=verify-full üzere PEM kodlu bir sunucu sertifikası sağla. Sertifika sağlamıyorsanız, bağlantı JVM varsayılan güven deposunu kullanır sslmode=verify-full .

Yapılandırma seçenekleri

Sink, tanınmayan seçeneklerde hata verir, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Sink yapılandırma seçenekleri, ortak seçenekler ve her bağlantı yöntemi için seçenekler için PostgreSQL sink seçeneklerine bakınız.

Veri türü eşlemeleri

Alıcı, mevcut bir Lakebase ya da harici PostgreSQL tablosuna yazmadan önce her DataFrame sütununun karşılık gelen hedef sütunla uyumlu olup olmadığını denetler.

Aşağıdaki tablo, Databricks Runtime 18 LTS ve üzerindeki türleri içerir:

Spark türü Otomatik olarak oluşturulan Lakebase tablo tipi Mevcut PostgreSQL tablolarındaki uyumlu türler
ByteType, ShortType smallint smallint
IntegerType integer integer
LongType bigint bigint
FloatType real real
DoubleType double precision double precision
DecimalType numeric numeric
StringType text varchar, text
VarcharType(n) varchar(n) varchar, text
CharType(n) char(n) char
BinaryType bytea bytea
BooleanType boolean boolean
TimestampType timestamptz timestamptz
TimestampNTZType timestamp timestamp
DateType date date
ArrayType, MapType, StructType, VariantType, NullType jsonb json, jsonb

Aşağıdaki tablo, Databricks Runtime 19 ve üzerindeki türleri içerir:

Spark türü Otomatik olarak oluşturulan Lakebase tablo tipi Mevcut PostgreSQL tablolarındaki uyumlu türler
DayTimeIntervalType, YearMonthIntervalType interval interval

Yukarı davranış

Bu seçenek upsertkey , hedef tablonun birincil anahtar sütunlarını tanımlar. Mevcut bir tablo için, içindeki upsertkey sütunlar tablonun ana anahtarıyla tam olarak eşleşmelidir. Seçeneği belirtmezseniz, sink birincil anahtarı tablodan okur. Sink'in oluşturduğu bir Lakebase tablosu için, upsertkey birincil anahtarı tanımlar. Bu seçeneği belirtmezseniz, sink tabloyu birincil anahtar olmadan oluşturur.

Hedef tabloda birincil anahtar olduğunda, sink PostgreSQL'in INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... sözdizimi ile upsert yapar. Hedef tabloda birincil anahtar olmadığında, sink insert işlemleri yapar. Bir sorgunun çıkış modu bu davranışı etkilemez.

Tüm birincil anahtar sütunları DataFrame'de bulunmalı ve sayısal veya dizi türleri gibi karşılaştırılabilir türleri kullanmalıdır.

Performans ayarı

Toplu işleme ve geri basınç

Koşullardan herhangi biri karşılandığında temizleme tetiklenir:

  • Arabellek, varsayılan değeri batchsize olan 1000 satıra ulaşır.
  • Arabellek yaşı, varsayılan değeri batchinterval olan 100 milliseconds değerini aşıyor.

Veritabanı gelen veri hızına ayak uyduramadığında, alıcı geri basıncı kaynağa doğru üst akışa iletir.

Gecikme süresi ve aktarım hızı kılavuzu:

  • Gerçek zamanlı moddaki düşük gecikmeli iş yüklerinde, boşaltma öncesindeki azami sürenin daha kısa olmasını sağlamak için batchinterval değerini azaltın. Kavramlar için Gerçek zamanlı mod kavramlarına , kod örneği için ise Gerçek zamanlı mod örneklerine bakınız.
  • Yüksek işlem hacimli iş yükleri için, işlem başına ek yükü azaltmak amacıyla batchsize değerini artırın.

Bağlantı davranışı

Sink, executorlarda bağlantı havuzu kullanır. Varsayılan olarak, her görev bir veritabanı bağlantısı kullanır.

Databricks, her bağlantı için varsayılan görev değerini 1 kullanmanızı önerir. Her bağlantı için görev sayısını artırırsanız, bağlantı çakışmalarına neden olabilir ve yüksek aktarım hızı bağlantıları için gecikme sürelerini artırabilirsiniz.

Görevlerin bağlantılara oranını yapılandırmak için Spark yapılandırmasını ayarlayın spark.databricks.sql.streaming.jdbc.tasksPerConnection . Hedef veritabanında düşük bir bağlantı sınırı varsa, shuffle bölüm sayısını azaltın veya spark.databricks.sql.streaming.jdbc.tasksPerConnection değerini artırın.

Sink, bağlantı hataları, kilitlenmeler ve oran sınırlaması dahil geçici JDBC hatalarını otomatik olarak yeniden dener. Havuz tüm yeniden denemeleri tüketirse sorgu başarısız olur.

Desteklenen tetikleyiciler ve çıkış modları

Triggers

Bu tablo, klasik ve sunucusuz hesaplamada Yapılandırılmış Akış tetikleyici türlerini destekleri gösterir:

Trigger Klasik işlem Sunucusuz işlem (not defterleri ve işler)
RealTime Yes Hayır
ProcessingTime Yes Hayır
AvailableNow Yes Yes
Once Yes. Deprecated. AvailableNow adresini kullanın. Yes. Deprecated. AvailableNow adresini kullanın.

Çıkış modları

Bu tabloda Yapılandırılmış Akış çıkış modları için destek gösterilmektedir:

Çıkış modu Supported
update Yes
append Yes. Davranış, update ile aynıdır. Hedef tabloda birincil anahtar varsa sorgu upsert işlemi yapar, aksi takdirde ekleme yapar. Upsert davranışına bakınız.
complete Hayır

Sınırlama

  • Unity Kataloğu bağlantısı üzerinden bağlanan harici bir PostgreSQL veritabanı için, hedef tablo zaten var olmalıdır. Sink, eksik tabloları yalnızca Lakebase'de otomatik olarak oluşturur.
  • Göl akışı boru hatları desteklenmemektedir.