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.
API'den veri alma, bir dosyadan veya veritabanından okumak yerine, genellikle sayfalandırılmış JSON biçiminde bir web hizmetinden HTTP üzerinden veri çekmek demektir. Dosyalar veya mesaj veri yolunun aksine, yerleşik genel bir API kaynağı yoktur, bu yüzden kimlik doğrulama, sayfa belirleme ve hız sınırlarını kendiniz halledersiniz. Lakeflow boru hatları, herhangi bir API'den veri alma için üç modeli destekler. Hangisinin uyması, ses ve yenileme ihtiyaçlarınıza bağlı.
Important
Herhangi bir özel API alım kodu yazmadan önce, kaynağınız için yönetilen bir bağlayıcının zaten var olup olmadığını kontrol edin. Lakeflow Connect, Salesforce, Workday, ServiceNow ve Google Analytics gibi birçok yaygın yazılım hizmet olarak (SaaS) API'si için yerleşik bağlayıcılar sunar ve ayrıca artan bir ortak bağlantı seti de var. Bir bağlayıcı veri kaynağınızı destekliyorsa, kimlik doğrulama, sayfalama ve artımlı veri çıkarma işlemlerini sizin yerinize yönetir; ayrıca bu, neredeyse her zaman elle oluşturulmuş bir veri alma sürecinden daha az çaba gerektirir. Bkz. Lakeflow Connect'te yönetilen bağlayıcılar. Aşağıdaki kalıpları yalnızca hiçbir konnektör uymadığında kullanın.
Prerequisites
- Bir boru hattı. Bir tane oluşturmak için Lakeflow pipelines eğitimlerine bakınız.
- API kimlik bilgileri, örneğin token veya anahtar, Azure Databricks gizlisi olarak saklanır. Pipeline kaynak kodunda kimlik bilgilerini asla sert kodlamayın. Bkz. Gizli yönetim.
- Pipeline hesaplamanızdan API uç noktasına ağ erişimi.
- Akış tablolarına ve somutlaştırılmış görünümlere, ayrıca bu kalıpların ürettiği veri kümesi türlerine aşina olma. Akış tabloları ve Maddeleştirilmiş görünümlere bakınız.
Desen seçme
Pipeline'larda yerleşik genel bir REST API kaynağı yoktur; bu nedenle, herhangi bir API'den veri çekerken, veri hacmine ve veriyi ne sıklıkta içeri aktardığınıza bağlı olarak üç yaklaşımdan birini seçin:
| Desen | Şu durumlarda kullanın: |
|---|---|
| Somutlaştırılmış görünüm olarak periyodik çekme işlemleri | Yükler küçük ve orta ölçekli olup her boru hattı çalışmasında bir kez çekilir; örneğin referans verisi, günlük FX oranları veya sayfalanmış ama sınırlanabilir bir API. |
| Python Veri Kaynak API'si | Yüksek hacimli veya akış tabanlı bir API’yi, denetim noktalarıyla kaydedilen ilerlemeyle birlikte artımlı olarak yoklamanız gerekir; böylece yeniden başlatma durumunda her şey baştan yeniden okunmaz. |
| Auto Loader ile ayrıştırılmış veri alımı | API’ye özgü kendine has davranışları dönüşüm mantığınızdan ayırmak ve dosyaların tam bir kez izlenmesini ek çaba harcamadan elde etmek istersiniz. |
Desen 1: Somutlaştırılmış bir görünüm olarak periyodik çekimler
Her boru hattı çalıştırma başına bir kez çekilen küçük-orta ölçekli yükler için, API'yi çağıran ve bir Spark DataFrame döndüren bir Python fonksiyonu yazın. Veri seti maddeleşmiş bir görünüm olduğundan, boru hattı her güncellemede işlevi tam ve aynı şekilde yeniden çalıştırır.
Aşağıdaki adımlar, belirli aralıklarla veri çekerek somutlaştırılmış görünümün nasıl oluşturulacağını gösterir:
API token'ını bir gizli içinde sakla, sonra pipeline ayarlarında bir Spark yapılandırma özelliğine eşlendirin ki pipeline kodu onu okuyabilsin. Bu özelliği
spark_confpipeline'ın küme konfigürasyonunun bloğuna ekleyin:{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }Bir sonraki adımda kod, bu değeri
spark.conf.get("api.token")ile okur. Boru hattı ayarlarında gizli bilgilerin yapılandırılması hakkında daha fazla bilgi için şu bölüme bakın: Boru hattındaki gizli bilgilerle depolama kimlik bilgilerine güvenli erişim.API'yi çağıran ve cevabı DataFrame olarak döndüren bir maddeleştirilmiş görünüm tanımlayın:
import requests from pyspark import pipelines as dp from pyspark.sql import Row @dp.materialized_view( name="exchange_rates_bronze", comment="Daily FX rates pulled from a public REST API", ) def exchange_rates_bronze(): resp = requests.get( "https://api.example.com/v1/rates", params={"base": "USD"}, headers={"Authorization": f"Bearer {spark.conf.get('api.token')}"}, timeout=30, ) resp.raise_for_status() rates = resp.json()["rates"] rows = [Row(currency=k, rate=float(v), as_of_date=resp.json()["date"]) for k, v in rates.items()] return spark.createDataFrame(rows)Fonksiyon içinde, sayfalar üzerinde döngü kurup sonuçları birleştirerek DataFrame'i döndürmeden önce sayfalamayı yönetin:
import requests from pyspark import pipelines as dp from pyspark.sql import Row @dp.materialized_view( name="customers_bronze", comment="Customers pulled from a paginated REST API", ) def customers_bronze(): token = spark.conf.get("api.token") rows = [] url = "https://api.example.com/v1/customers" while url: # follow the API's next-page cursor until exhausted resp = requests.get( url, headers={"Authorization": f"Bearer {token}"}, timeout=30, ) resp.raise_for_status() payload = resp.json() rows.extend(Row(**record) for record in payload["data"]) url = payload.get("next") # None on the last page return spark.createDataFrame(rows)Dayanıklılık talebi etrafında yeniden deneme ve geri çekilme mantığı ekleyin.
Bu desen, her boru hattı güncellemesinde tam API yanıtını tekrar okur, bu yüzden sadece yük sınırlı olduğunda kullanın. Artımlı okumalar için 2 numaralı deseni kullanın.
Desen 2: Python Veri Kaynağı API'si ile yüksek hacimli veya akışlı API'ler
API'ler için, offset takibi ile kademeli olarak anket yapmanız gerekiyor, Spark'ın Python Data Source API'sini kullanarak özel bir veri kaynağı uygulayın. Bu, denetim noktasıyla kaydedilen ilerleme ve artımlı okumalar dâhil olmak üzere size uygun akış semantiği sağlar; böylece yeniden başlatıldığında API’nin tamamı yeniden çekilmek yerine son ofsetten devam edilir.
Aşağıdaki adımlar, özel bir veri kaynağından nasıl alım yapacağınızı gösterir:
API'yi çağıran ve okuma ofsetini izleyen bir
DataSourceveDataSourceStreamReaderuygulayın. Özel veri kaynağı oluşturma hakkında detaylar için PySpark özel veri kaynakları sayfasına bakınız.Veri kaynağını kaydedin ki boru hattı format adıyla referans versin:
spark.dataSource.register(MyApiDataSource)Kayıtlı kaynaktan bir akış tablosunda okuyun:
from pyspark import pipelines as dp @dp.table(name="events_bronze") def events_bronze(): return spark.readStream.format("my_api_source").load()
Desen 3: Veri alımını zamanlanmış bir iş ve Auto Loader ile ayrıştırma
Yaygın bir üretim modeli, API çağrısını boru hattından ayırmaktır. Zamanlanmış bir iş, ham API yanıtlarını Unity Kataloğu hacminde dosya olarak getirir ve pipeline bunları Auto Loader ile alır. Bu, sayfalama ve oran sınırlamaları gibi API'ye özgü ayrıntıları bildirimsel dönüştürme mantığınızdan yalıtır ve size Auto Loader'ın tam olarak bir kez dosya izleme özelliğini ek bir çaba gerektirmeden sunar.
Aşağıdaki adımlar, alma işleminin zamanlanmış bir iş ile bağlantısını nasıl keseceğinizi gösterir:
API'yi çağıran ve ham JSON yanıtlarını bir Unity Kataloğu hacmine yazan bir not defteri veya betik yaz. API kimlik bilgilerini bir gizlikten okuyun. Bkz. Gizli yönetim.
import requests, json, time token = dbutils.secrets.get(scope="<scope-name>", key="<secret-name>") volume_path = "/Volumes/main/raw/landing/api_events" resp = requests.get( "https://api.example.com/v1/events", headers={"Authorization": f"Bearer {token}"}, timeout=30, ) resp.raise_for_status() # One file per run; the pipeline's Auto Loader tracks which files it has ingested. with open(f"{volume_path}/events_{int(time.time())}.json", "w") as f: json.dump(resp.json()["data"], f)Not defterini veya betiği Lakeflow Jobs ile otomatik olarak çalışacak şekilde zamanlayın. Bakınız Lakeflow İşleri.
Boru hattınızda, Otomatik Yükleyici ile indirilmiş dosyaları okuyan bir akış tablosu tanımlayın:
from pyspark import pipelines as dp @dp.table(name="api_events_bronze") def api_events_bronze(): return ( spark.readStream.format("cloudFiles") .option("cloudFiles.format", "json") .load("/Volumes/main/raw/landing/api_events") )
Auto Loader ile güvenilir dosya girişi hakkında daha fazla bilgi için Bulut nesne depolamasından dosyaları yükleme ve Otomatik Yükleyici Nedir? bölümlerine bakabilirsiniz.
API alımı için en iyi uygulamalar
- Sırları kaynak kodundan uzak tutun. API token'larını ve anahtarlarını Azure Databricks gizli kapsamlarında depolayın ve çalışma zamanında okuyun. Bkz. Gizli yönetim.
- Yanıtları erken onaylayın. Hatalı biçimlendirilmiş API yanıtlarını alt süreçlere iletilmeden önce yakalamak için, içe aktarılan satırlara beklentiler ekleyin.
- Sayfa ve hız sınırlarını yönetin. Sayfalar üzerinde döngü kur ve artan aralıklarla yeniden deneme ekle; böylece geçici bir hata tüm güncellemenin başarısız olmasına neden olmasın.