Apache Spark için Azure Veri Gezgini Bağlayıcısı

Apache Spark , büyük ölçekli veri işlemeye yönelik birleşik bir analiz altyapısıdır. Azure Veri Gezgini büyük miktarda veri üzerinde gerçek zamanlı analiz yapmaya yönelik hızlı ve tam olarak yönetilen bir veri analizi hizmetidir.

Spark için Kusto bağlayıcısı, herhangi bir Spark kümesinde çalışabilen açık kaynak bir projedir. Verileri Azure Veri Gezgini ve Spark kümeleri arasında taşımak için veri kaynağı ve veri havuzu uygular. Azure Veri Gezgini ve Apache Spark kullanarak, veri odaklı senaryoları hedefleyen hızlı ve ölçeklenebilir uygulamalar oluşturabilirsiniz. Örneğin makine öğrenmesi (ML), Ayıklama-Dönüştürme-Yükleme (ETL) ve Log Analytics. Bağlayıcı ile Azure Veri Gezgini, , writeve readgibi writeStreamstandart Spark kaynağı ve havuz işlemleri için geçerli bir veri deposu haline gelir.

Kuyruklu veri alma veya akışla veri alma yöntemiyle Azure Veri Gezgini'ne yazabilirsiniz. Azure Veri Gezgini'nden okuma, Azure Veri Gezgini'ndeki verileri filtreleyerek aktarılan verilerin hacmini azaltan kolon kırpma ve koşul itme özelliklerini destekler.

Bu makalede, Azure Veri Gezgini Spark bağlayıcısının nasıl yükleneceği ve yapılandırıldığı ve Verilerin Azure Veri Gezgini ile Apache Spark kümeleri arasında nasıl taşındığı açıklanır.

Not

Bu makaledeki örneklerden bazıları Azure Databricks Spark kümesine başvursa da, Azure Veri Gezgini Spark bağlayıcısı Databricks'e veya başka bir Spark dağıtımına doğrudan bağımlılık almaz.

Önkoşullar

İpucu

Spark 2.3.x sürümleri de desteklenir, ancak pom.xmliçindeki bazı bağımlılıkları değiştirmeniz gerekebilir.

Spark bağlayıcısı oluşturma

2.3.0 sürümünden başlayarak spark-kusto-connector: kusto-spark_3.0_2.12 yerine Spark 3.x ve Scala 2.12'yi hedefleyen yeni yapıt kimliklerini kullanıma sunacağız.

Not

2.5.1 öncesi sürümler mevcut bir tabloya almak için artık çalışmıyor, lütfen daha sonraki bir sürüme güncelleştirin. Bu adım isteğe bağlıdır. Maven gibi önceden oluşturulmuş kitaplıklar kullanıyorsanız bkz . Spark kümesi kurulumu.

Önkoşullar

  1. Spark Bağlayıcısı oluşturmak için bu kaynağa bakın.

  2. Maven proje tanımlarını kullanan Scala/Java uygulamaları için uygulamanızı en son yapıtla bağlayın. Maven Central'de en son artefaktı bulun.

    For more information, see [https://mvnrepository.com/artifact/com.microsoft.azure.kusto/kusto-spark_3.0_2.12](https://mvnrepository.com/artifact/com.microsoft.azure.kusto/kusto-spark_3.0_2.12).
    
    
  3. Önceden oluşturulmuş kitaplıklar kullanmıyorsanız, aşağıdaki Kusto Java SDK kitaplıkları da dahil olmak üzere bağımlılıklarda listelenen kitaplıkları yüklemeniz gerekir. Yüklenecek doğru sürümü bulmak için ilgili sürümün pom'una bakın:

    1. Jar oluşturmak ve tüm testleri çalıştırmak için:

      mvn clean package -DskipTests
      
    2. Jar oluşturmak için tüm testleri çalıştırın ve jar dosyasını yerel Maven deponuza yükleyin:

      mvn clean install -DskipTests
      

Daha fazla bilgi için bkz . bağlayıcı kullanımı.

Spark kümesi kurulumu

Not

Aşağıdaki adımları gerçekleştirirken en son Kusto Spark bağlayıcı sürümünü kullanmanız önerilir.

  1. Azure Databricks kümesi Spark 3.0.1 ve Scala 2.12'yi temel alarak aşağıdaki Spark kümesi ayarlarını yapılandırın:

    Databricks kümesi ayarları.

  2. Maven'dan en son spark-kusto-connector kitaplığını yükleyin:

    Kitaplıkları içeri aktar. Spark-Kusto-Connector'ı seçin.

  3. Tüm gerekli kitaplıkların yüklendiğini doğrulayın:

    Yüklü kitaplıkları doğrulayın.

  4. JAR dosyası kullanarak yükleme için diğer bağımlılıkların yüklendiğini doğrulayın:

    Bağımlılıkları ekleyin.

Kimlik Doğrulaması

Kusto Spark bağlayıcısı, aşağıdaki yöntemlerden birini kullanarak Microsoft Entra Id ile kimlik doğrulaması yapmanıza olanak tanır:

Microsoft Entra uygulaması kimlik doğrulaması

Microsoft Entra uygulaması kimlik doğrulaması en basit ve en yaygın kimlik doğrulama yöntemidir ve Kusto Spark bağlayıcısı için önerilir.

  1. Azure CLI aracılığıyla Azure aboneliğinizde oturum açın. Ardından tarayıcıda kimlik doğrulaması yapın.

    az login
    
  2. Sorumluyu barındırmak için aboneliği seçin. Bu adım, birden çok aboneliğiniz olduğunda gereklidir.

    az account set --subscription YOUR_SUBSCRIPTION_GUID
    
  3. Hizmet sorumlusunu oluşturun. Bu örnekte, hizmet sorumlusunun adı my-service-principal olarak adlandırılır.

    az ad sp create-for-rbac -n "my-service-principal" --role Contributor --scopes /subscriptions/{SubID}
    
  4. Döndürülen JSON verilerinden, gelecekte kullanmak üzere appId, password ve tenant değerini kopyalayın.

    {
      "appId": "00001111-aaaa-2222-bbbb-3333cccc4444",
      "displayName": "my-service-principal",
      "name": "my-service-principal",
      "password": "00001111-aaaa-2222-bbbb-3333cccc4444",
      "tenant": "00001111-aaaa-2222-bbbb-3333cccc4444"
    }
    

Microsoft Entra uygulamanızı ve hizmet sorumlunuzu oluşturdunuz.

Spark bağlayıcısı kimlik doğrulaması için aşağıdaki Entra uygulaması özelliklerini kullanır:

Özellikler Seçenek Dizesi Açıklama
KUSTO_AAD_APP_ID kustoAadAppId Microsoft Entra uygulaması (istemci) tanımlayıcısı.
KUSTO_AAD_AUTHORITY_ID kustoAadAuthorityID Microsoft Entra kimlik doğrulama yetkilisi. Microsoft Entra Directory (kiracı) kimliği. İsteğe bağlı - varsayılan olarak microsoft.com. Daha fazla bilgi için bkz . Microsoft Entra authority.
KUSTO_AAD_APP_SECRET kustoAadAppSecret İstemci için Microsoft Entra uygulama anahtarı.
KUSTO_ACCESS_TOKEN kustoAccessToken Eğer Kusto'ya erişim sağlayan bir accessToken'iniz zaten varsa, bu, kimlik doğrulaması için bağlayıcıya da geçirilebilir.

Not

Eski API sürümleri (2.0.0'dan küçük) şu adlandırmaya sahiptir: "kustoAADClientID", "kustoClientAADClientPassword", "kustoAADAuthorityID"

Kusto ayrıcalıkları

Gerçekleştirmek istediğiniz Spark işlemine göre kusto tarafında aşağıdaki ayrıcalıkları verin.

Spark uygulaması işlemi Ayrıcalıklar
Okuma - Tekli Mod Okuyucu
Okuma – Dağıtılmış Modu Zorla Okuyucu
Kuyruğa Alınmış Mod – CreateTableIfNotExist tablo oluşturma seçeneğiyle Yazma Yönetici
Yazma – Kuyruğa Alınmış Mod (FailIfNotExist tablo oluşturma seçeneğiyle) İçeri Alıcı
Yazma – İşlemsel Mod Yönetici

Asıl roller hakkında daha fazla bilgi için bkz . rol tabanlı erişim denetimi. Güvenlik rollerini yönetmek için bkz . güvenlik rolleri yönetimi.

Spark havuzu: Kusto'ya yazma

  1. Havuz parametrelerini ayarlayın:

    val KustoSparkTestAppId = dbutils.secrets.get(scope = "KustoDemos", key = "KustoSparkTestAppId")
    val KustoSparkTestAppKey = dbutils.secrets.get(scope = "KustoDemos", key = "KustoSparkTestAppKey")
    
    val appId = KustoSparkTestAppId
    val appKey = KustoSparkTestAppKey
    val authorityId = "72f988bf-86f1-41af-91ab-2d7cd011db47" // Optional - defaults to microsoft.com
    val cluster = "Sparktest.eastus2"
    val database = "TestDb"
    val table = "StringAndIntTable"
    
  2. Spark DataFrame'i Kusto kümesine toplu olarak yazın:

    import com.microsoft.kusto.spark.datasink.KustoSinkOptions
    import org.apache.spark.sql.{SaveMode, SparkSession}
    
    df.write
      .format("com.microsoft.kusto.spark.datasource")
      .option(KustoSinkOptions.KUSTO_CLUSTER, cluster)
      .option(KustoSinkOptions.KUSTO_DATABASE, database)
      .option(KustoSinkOptions.KUSTO_TABLE, "Demo3_spark")
      .option(KustoSinkOptions.KUSTO_AAD_APP_ID, appId)
      .option(KustoSinkOptions.KUSTO_AAD_APP_SECRET, appKey)
      .option(KustoSinkOptions.KUSTO_AAD_AUTHORITY_ID, authorityId)
      .option(KustoSinkOptions.KUSTO_TABLE_CREATE_OPTIONS, "CreateIfNotExist")
      .mode(SaveMode.Append)
      .save()  
    

    Basitleştirilmiş söz dizimini de kullanabilirsiniz:

    import com.microsoft.kusto.spark.datasink.SparkIngestionProperties
    import com.microsoft.kusto.spark.sql.extension.SparkExtension._
    
    // Optional, for any extra options:
    val conf: Map[String, String] = Map()
    
    val sparkIngestionProperties = Some(new SparkIngestionProperties()) // Optional, use None if not needed
    df.write.kusto(cluster, database, table, conf, sparkIngestionProperties)
    
  3. Akış verileri yazma:

    import org.apache.spark.sql.streaming.Trigger
    import java.util.concurrent.TimeUnit
    import java.util.concurrent.TimeUnit
    import org.apache.spark.sql.streaming.Trigger
    
    // Set up a checkpoint and disable codeGen. 
    spark.conf.set("spark.sql.streaming.checkpointLocation", "/FileStore/temp/checkpoint")
    
    // As an alternative to adding .option by .option, you can provide a map:
    val conf: Map[String, String] = Map(
      KustoSinkOptions.KUSTO_CLUSTER -> cluster,
      KustoSinkOptions.KUSTO_TABLE -> table,
      KustoSinkOptions.KUSTO_DATABASE -> database,
      KustoSourceOptions.KUSTO_ACCESS_TOKEN -> accessToken)
    
    // Write to a Kusto table from a streaming source
    val kustoQ = df
      .writeStream
      .format("com.microsoft.kusto.spark.datasink.KustoSinkProvider")
      .options(conf)
      .trigger(Trigger.ProcessingTime(TimeUnit.SECONDS.toMillis(10))) // Sync this with the ingestionBatching policy of the database
      .start()
    

Spark kaynağı: Kusto'dan okuma

  1. Az miktarda veri okurken veri sorgusunu tanımlayın:

    import com.microsoft.kusto.spark.datasource.KustoSourceOptions
    import org.apache.spark.SparkConf
    import org.apache.spark.sql._
    import com.microsoft.azure.kusto.data.ClientRequestProperties
    
    val query = s"$table | where (ColB % 1000 == 0) | distinct ColA"
    val conf: Map[String, String] = Map(
          KustoSourceOptions.KUSTO_AAD_APP_ID -> appId,
          KustoSourceOptions.KUSTO_AAD_APP_SECRET -> appKey
        )
    
    val df = spark.read.format("com.microsoft.kusto.spark.datasource").
      options(conf).
      option(KustoSourceOptions.KUSTO_QUERY, query).
      option(KustoSourceOptions.KUSTO_DATABASE, database).
      option(KustoSourceOptions.KUSTO_CLUSTER, cluster).
      load()
    
    // Simplified syntax flavor
    import com.microsoft.kusto.spark.sql.extension.SparkExtension._
    
    val cpr: Option[ClientRequestProperties] = None // Optional
    val df2 = spark.read.kusto(cluster, database, query, conf, cpr)
    display(df2)
    
  2. İsteğe bağlı: Geçici blob depolama (Kusto değil) sağlarsanız, bloblar çağıranın sorumluluğu altında oluşturulur. Buna depolama sağlama, erişim anahtarlarını döndürme ve geçici yapıtları silme dahildir. KustoBlobStorageUtils modülü, hesap ve kapsayıcı koordinatlarına ve hesap kimlik bilgilerine göre blobları silmeye yönelik yardımcı işlevler ya da yazma, okuma ve liste izinlerine sahip tam bir SAS URL'si içerir. Karşılık gelen RDD artık gerekli olmadığında, her işlem, geçici blob artefaktlarını ayrı bir dizinde depolar. Bu dizin, Spark Sürücüsü düğümünde bildirilen okuma işlemi bilgi günlüklerinin bir bölümü olarak kaydedilir.

    // Use either container/account-key/account name, or container SaS
    val container = dbutils.secrets.get(scope = "KustoDemos", key = "blobContainer")
    val storageAccountKey = dbutils.secrets.get(scope = "KustoDemos", key = "blobStorageAccountKey")
    val storageAccountName = dbutils.secrets.get(scope = "KustoDemos", key = "blobStorageAccountName")
    // val storageSas = dbutils.secrets.get(scope = "KustoDemos", key = "blobStorageSasUrl")
    

    Yukarıdaki örnekte Key Vault'a bağlayıcı arabirimi kullanılarak erişilenmez; Databricks gizli dizilerini kullanmanın daha basit bir yöntemi kullanılır.

  3. Kusto'dan okuyun.

    • Geçici blob depolamayı sağlarsanız Kusto'dan şu şekilde okuyun:

       val conf3 = Map(
            KustoSourceOptions.KUSTO_AAD_APP_ID -> appId,
            KustoSourceOptions.KUSTO_AAD_APP_SECRET -> appKey
            KustoSourceOptions.KUSTO_BLOB_STORAGE_SAS_URL -> storageSas)
      val df2 = spark.read.kusto(cluster, database, "ReallyBigTable", conf3)
      
      val dfFiltered = df2
        .where(df2.col("ColA").startsWith("row-2"))
        .filter("ColB > 12")
        .filter("ColB <= 21")
        .select("ColA")
      
      display(dfFiltered)
      
    • Kusto, geçici blob depolamayı sağlıyorsa, aşağıdaki şekilde Kusto'dan okuyun:

      val conf3 = Map(
        KustoSourceOptions.KUSTO_AAD_CLIENT_ID -> appId,
        KustoSourceOptions.KUSTO_AAD_CLIENT_PASSWORD -> appKey)
      val df2 = spark.read.kusto(cluster, database, "ReallyBigTable", conf3)
      
      val dfFiltered = df2
        .where(df2.col("ColA").startsWith("row-2"))
        .filter("ColB > 12")
        .filter("ColB <= 21")
        .select("ColA")
      
      display(dfFiltered)