SQL veritabanları için Spark bağlayıcısı

SQL veritabanları için Spark bağlantısı, SQL Server, Azure SQL veritabanları ve Fabric'teki SQL veritabanlarından okuma ve yazma imkanı veren yüksek performanslı bir kütüphanedir. Bağlayıcı aşağıdaki özellikleri sunar:

  • Azure SQL Veritabanı, Azure SQL Yönetilen Örneği, Azure VM'de SQL Server ve Fabric'te SQL veritabanlarında büyük yazma ve okuma işlemleri yapmak için Spark kullanın.
  • Tablo veya görünüm kullandığınızda, bağlayıcı SQL altyapısı düzeyinde ayarlanan güvenlik modellerini destekler. Bu modeller nesne düzeyinde güvenlik (OLS), satır düzeyi güvenlik (RLS) ve sütun düzeyinde güvenlik (CLS) içerir.

Bağlayıcı, Fabric çalışma zamanına öncesinden yüklendiği için ayrı olarak yüklemeniz gerekmez.

Authentication

Microsoft Entra doğrulaması Fabric ile entegre edilmiştir.

  • Doku çalışma alanında oturum açtığınızda kimlik bilgileriniz kimlik doğrulaması ve yetkilendirme için otomatik olarak SQL altyapısına geçirilir.
  • SQL veritabanı altyapınızda Microsoft Entra Kimliğinin etkinleştirilmesini ve yapılandırılmasını gerektirir.
  • Microsoft Entra Id ayarlandıysa Spark kodunuzda ek yapılandırma gerekmez. Kimlik bilgileri otomatik olarak eşlenir.

SQL kimlik doğrulama yöntemini (SQL kullanıcı adı ve parola belirterek) veya hizmet sorumlusunu da kullanabilirsiniz (uygulama tabanlı kimlik doğrulaması için Azure erişim belirteci sağlayarak).

Permissions

Spark bağlayıcısını kullanmak için kimliğinizin (ister kullanıcı ister uygulama olsun) hedef SQL altyapısı için gerekli veritabanı izinlerine sahip olması gerekir. Bu izinler, tablolardan ve görünümlerden okumak veya bu görünümlere yazmak için gereklidir.

Azure SQL Veritabanı, Azure SQL Yönetilen Örneği ve Azure VM üzerinde SQL Server için:

  • Genellikle, işlemi çalıştıran kimlik, db_datawriter ve db_datareader izinlerine ihtiyaç duyar ve tam denetim için isteğe bağlı olarak db_owner gereklidir.

Fabric'te bir SQL veritabanı için:

  • Kimlik genellikle db_datawriter ve db_datareader izinlerine ihtiyaç duyar ve isteğe bağlı olarak db_owner izinlerine ihtiyaç duyabilir.
  • Kimlik ayrıca öğe seviyesinde Fabric'teki SQL veritabanında en az okuma iznine sahip olmalıdır.

Uyarı

Hizmet ilkesi kullanıyorsanız, kullanıcı kimliğine bürünme etkinleştirildiyse bir uygulama olarak (kullanıcı bağlamı olmadan) veya kullanıcı olarak çalışabilir. Hizmet sorumlusu, gerçekleştirmek istediğiniz işlemler için gerekli veritabanı izinlerine sahip olmalıdır.

Kullanım ve kod örnekleri

Bu bölümde, SQL veritabanları için Spark bağlayıcısının nasıl etkili bir şekilde kullanılacağını gösteren kod örnekleri sağlıyoruz. Bu örnekler, SQL tablolarından okuma ve sql tablolarına yazma ve bağlayıcı seçeneklerini yapılandırma gibi çeşitli senaryoları kapsar.

Uyarı

Toplu yazmadan önce, gelen tüm Spark verileri hedef SQL veri tiplerine uymalıdır. Bir tablonun üzerine yazdığınızda veya oluşturduğunuzda, bağlayıcı Spark TimestampType ve TimestampNTZType değerlerini SQL datetime yerine SQL datetime2 ile eşler. Kıvılcım zaman damgası türleri altı haneli kesirlik saniye hassasiyetini desteklerken, SQL datetime üç haneyi destekler ve bu da uyumsuzluğa yol açabilir.

Desteklenen seçenekler

Gerekli en düşük seçenek url olarak "jdbc:sqlserver://<server>:<port>;database=<database>;" veya spark.mssql.connector.default.url olarak ayarlamaktır.

  • url sağlandığında:

    • Her zaman ilk tercih olarak kullanın url .
    • Ayarlı değilse spark.mssql.connector.default.url bağlayıcı bunu ayarlar ve gelecekteki kullanımlar için yeniden kullanabilir.
  • url sağlanmadığında:

    • Eğer spark.mssql.connector.default.url ayarlanırsa, bağlayıcı spark yapılandırması değerini kullanır.
    • spark.mssql.connector.default.url ayarlanmamışsa, gerekli ayrıntılar mevcut olmadığından bir hata oluşur.

Bu bağlayıcı burada tanımlanan seçenekleri destekler: SQL DataSource JDBC Seçenekleri

Bağlayıcı aşağıdaki seçenekleri de destekler:

Seçenek Varsayılan değer Description
reliabilityLevel en iyi çaba Ekleme işlemlerinin güvenilirliğini denetler. Olası değerler: BEST_EFFORT (varsayılan, en hızlı, yürütücü yeniden başlatılırsa yinelenen satırlara neden olabilir) NO_DUPLICATES (daha yavaş, yürütücü yeniden başlatıldığında bile yinelenen satırların eklenmemesini sağlar). Yinelemelere ve performans gereksinimlerine olan toleransınızı temel alarak seçin.
isolationLevel "READ_COMMITTED" SQL işlemleri için işlem yalıtım düzeyini ayarlar. Olası değerler: READ_COMMITTED (varsayılan olarak, kaydedilmemiş verilerin okunmasını engeller), READ_UNCOMMITTED, REPEATABLE_READ, SNAPSHOT, SERIALIZABLE. Daha yüksek yalıtım düzeyleri eşzamanlılığı azaltabilir ancak veri tutarlılığını iyileştirebilir.
tableLock yanlış Ekleme işlemleri sırasında SQL Server TABLOCK tablo düzeyinde kilit ipucunun kullanılıp kullanılmayacağını denetler. Olası değerler: true (toplu yazma performansını geliştirebilen TABLOCK'yi etkinleştirir), false (varsayılan olarak TABLOCK'yi kullanmaz). ayarı true , büyük eklemeler için aktarım hızını artırabilir ancak tablodaki diğer işlemler için eşzamanlılığı azaltabilir.
schemaCheckEnabled doğru Spark'ınız DataFrame ve SQL tablonuz arasında katı şema doğrulamasının zorunlu kılınıp uygulanmayacağını denetler. Olası değerler: true (varsayılan olarak, katı şema eşleştirmesi uygular), false (daha fazla esneklik sağlar ve bazı şema denetimlerini atlayabilir). ayarı false şema uyuşmazlıklarına yardımcı olabilir, ancak yapılar önemli ölçüde farklıysa beklenmeyen sonuçlara yol açabilir.

Bazı Toplu API seçenekleriDataFrame üzerinde seçenekler olarak ayarlanabilir ve yazma sırasında toplu kopyalama API'lerine geçirilir.

Yazma ve okuma örneği

Aşağıdaki kod, bu işlemleri göstermek için otomatik Microsoft Entra ID doğrulama kullanır:

  • Bir DataFrame Yaz: df.write.option("...", "...").mssql("<schema>.<table>").
  • Bir tabloyu okuyun: spark.read.option("...", "...").mssql("<schema>.<table>").
  • Özel bir sorgu çalıştırın: spark.read.option("...", "...").option("query", "<your-custom-query>").mssql().

Tip

Veriler gösterim amacıyla satır içinde oluşturulur. Üretim senaryosunda genellikle mevcut bir kaynaktan verileri okur veya daha karmaşık DataFramebir oluşturursunuz.

import com.microsoft.sqlserver.jdbc.spark
url = "jdbc:sqlserver://<server>:<port>;database=<database>;"
row_data = [("Alice", 1),("Bob", 2),("Charlie", 3)]
column_header = ["Name", "Age"]
df = spark.createDataFrame(row_data, column_header)
df.write.mode("overwrite").option("url", url).mssql("dbo.publicExample")
spark.read.option("url", url).mssql("dbo.publicExample").show()
spark.read.option("url", url).option("query", "SELECT * FROM dbo.publicExample WHERE Age = 3").mssql().show() # Read with a custom query

url = "jdbc:sqlserver://<server>:<port>;database=<database2>;" # different database
df.write.mode("overwrite").option("url", url).mssql("dbo.tableInDatabase2") # default url is updated
spark.read.mssql("dbo.tableInDatabase2").show() # no url option specified and will use database2

Ayrıca SQL veritabanı altyapısından veri okurken sütunları seçebilir, filtre uygulayabilir ve diğer seçenekleri kullanabilirsiniz.

Kimlik doğrulama örnekleri

Aşağıdaki örneklerde hizmet sorumlusu (erişim belirteci) ve SQL kimlik doğrulaması gibi Microsoft Entra ID dışındaki kimlik doğrulama yöntemlerinin nasıl kullanılacağı gösterilmektedir.

Uyarı

Daha önce belirtildiği gibi, Doku çalışma alanında oturum açtığınızda Microsoft Entra Id kimlik doğrulaması otomatik olarak işlenir, bu nedenle bu yöntemleri yalnızca senaryonuz gerektiriyorsa kullanmanız gerekir.

import com.microsoft.sqlserver.jdbc.spark
url = "jdbc:sqlserver://<server>:<port>;database=<database>;"
row_data = [("Alice", 1),("Bob", 2),("Charlie", 3)]
column_header = ["Name", "Age"]
df = spark.createDataFrame(row_data, column_header)

from azure.identity import ClientSecretCredential
credential = ClientSecretCredential(tenant_id="", client_id="", client_secret="") # service principal app
scope = "https://database.windows.net/.default"
token = credential.get_token(scope).token

df.write.mode("overwrite").option("url", url).option("accesstoken", token).mssql("dbo.publicExample")
spark.read.option("accesstoken", token).mssql("dbo.publicExample").show()
spark.read.option("accesstoken", token).option("query", "SELECT * FROM dbo.publicExample WHERE Age = 3").mssql().show() # Read with a custom query

Desteklenen DataFrame kaydetme modları

Spark'tan SQL veritabanlarına veri yazarken çeşitli kaydetme modlarından birini seçebilirsiniz. Kaydetme modları, hedef tablo zaten mevcut olduğunda verilerin nasıl yazılması olduğunu denetler ve şemayı, verileri ve dizin oluşturmayı etkileyebilir. Bu modları anlamak beklenmeyen veri kaybını veya değişikliklerini önlemenize yardımcı olur.

Bu bağlayıcı burada tanımlanan seçenekleri destekler: Spark Save işlevleri

  • ErrorIfExists (varsayılan kaydetme modu): Hedef tablo varsa, yazma iptal edilir ve bir özel durum döndürülür. Aksi takdirde, verilerle yeni bir tablo oluşturulur.

  • Yoksay: Hedef tablo varsa, yazma isteği görmezden gelinir ve bir hata döndürülmez. Aksi takdirde, verilerle yeni bir tablo oluşturulur.

  • Üzerine yaz: Hedef tablo varsa, tablo bırakılır, yeniden oluşturulur ve yeni veriler eklenir.

    Uyarı

    Kullandığınızda overwrite, orijinal tablo şemasını (özellikle MSSQL'ye özel veri türlerini) ve tablo indekslerini kaybedersiniz. Şema, Spark DataFrame'inizden çıkarılan şema ile değiştirilir. Şemayı ve indeksleri kaybetmemek için .option("truncate", true) ekleyin.

  • Ekleme: Hedef tablo varsa, yeni veriler eklenir. Aksi takdirde, verilerle yeni bir tablo oluşturulur.

Troubleshoot

İşlem tamamlandığında Spark okuma işleminizin çıkışı hücrenin çıkış alanında görünür. Hatalar com.microsoft.sqlserver.jdbc.SQLServerException doğrudan SQL Server'dan gelir. Ayrıntılı hata bilgilerini Spark uygulama günlüklerinde bulabilirsiniz.

Toplu yazmalar, hedef SQL veri tiplerine uymak için gelen veri gerektirir. Veriler uymuyorsa, şu hata alabilirsiniz:

Caused by: com.microsoft.sqlserver.jdbc.SQLServerException: The service has encountered an error processing your request. Please try again. Error code 4815.

Örneğin, hedef SQL tablosu datetime kullandığında, 2025-01-01 10:30:00.123456 gibi gelen bir Spark TimestampType değeri, datetime öğesinin desteklediğinden daha fazla hassasiyete sahip olduğunda bu hata oluşabilir.

Hatayı çözmek için şu yaklaşımlardan birini kullanın:

  • Gelen veriyi hedef SQL veri tipine uyacak şekilde dök, kesik veya dönüştür. Örneğin, değeri üç haneli kesir-saniye hassasiyetine indirin: 2025-01-01 10:30:00.123000.
  • .option("truncate", false) ayarını yaparak bağlayıcının tablo şemasını yeniden oluşturmasına izin verin. Bağlayıcı, Spark zaman damgası tipini SQL datetime2ile eşler.