Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Az Azure Databricksben az Apache Spark által támogatott bármely adatforrásból betölthet adatokat folyamatok használatával. A folyamat adatkészleteit (táblázatait és nézeteit) bármely olyan lekérdezéshez definiálhatja, amely Spark DataFrame-et ad vissza, beleértve a Stream DataFrame-eket és a Pandas for Spark DataFrames-et. Az adatbetöltési feladatokhoz a Databricks a streamelési táblák használatát javasolja a legtöbb használati esetben. Az adatfolyam-táblák hasznosak adatok betöltésére felhőalapú objektumtárolóból Auto Loader segítségével vagy üzenettovábbító rendszerekből, például a Kafkából. A streamelési táblákról a betöltéshez használt elsődleges adatkészlettípusról további információt a Streamelési táblák című témakörben talál.
Nem minden adatforrás rendelkezik SQL-támogatással a betöltéshez. SQL- és Python-forrásokat kombinálhat ugyanabban a folyamatban, hogy szükség esetén Python-t használjon. A folyamatokkal alapértelmezés szerint nem csomagolt kódtárak kezeléséről további információt a folyamatok Python függőségeinek kezelése című témakörben talál. Az Azure Databricksbe történő betöltéssel kapcsolatos általános információkért lásd: Standard összekötők a Lakeflow Connectben.
Az alábbi példák néhány gyakori adatbetöltési mintát mutatnak be.
Azonosítsd az adatforrásokat és a kapcsolati útvonalat
Mielőtt pipeline kódot írsz, felleltározz minden helyet, ahonnan az adat származik. Minden forrásnál figyeld meg, hogyan teszi ki az adatait (fájlok, adatbázis, szoftverként szolgáltatás (SaaS) rendszer, API vagy folyam), milyen gyakran változik, és milyen hitelesítésekre és hálózati hozzáférésre van szüksége. A kapcsolódási mód gyakran meghatározza, hogy egy forrás alapvetően kötegelt vagy streamalapú, ezért ha ezt már korán helyesen határozzuk meg, elkerülhető a későbbi átdolgozás.
Sorolj minden forrást az alábbi csatlakozási útvonalak egyikébe. Az alábbi táblázat felsorolja az egyes csővezeték-források preferált mechanizmusait:
| Source | Kapcsolati útvonal |
|---|---|
| A felhőalapú objektumtárolóba (S3, Azure Data Lake Storage, GCS) érkező fájlok | A leggyakoribb kiindulópont. Használd az Auto Loadert (cloudFiles formátum), amely kezeli a fokozatos felfedezést, sémák következtetéseket és sémafejlődést. Lásd: Fájlok betöltése a felhőobjektum-tárolóból. |
| Adatbázisok és SaaS alkalmazások (Salesforce, SQL Server, PostgreSQL, Workday) | Használj egy Lakeflow Connect menedzselt csatlakozót, ahol az létezik a forrásodhoz. A menedzselt csatlakozók konfigurációvezéreltek, és helyetted kezelik a hitelesítést és az inkrementális vagy CDC kinyerést. Lásd: Felügyelt csatlakozók a Lakeflow Connectben. Ha nincs menedzselt csatlakozó a forrásodhoz, először vedd be közvetlenül vagy add le a válaszokat fájlként. Lásd: Adatok betöltése API-ból folyamatokban. |
| Üzenetsínek (Kafka, Kinesis, Azure Event Hubs, Pub/Sub) | Közvetlenül Strukturált Streaming forrásként olvassuk, mert ezek natív streaming források. Lásd: Adat betöltése üzenetbuszról. |
| Egyéb Delta táblák vagy Unity katalógus eszközök, beleértve más csővezetékek vagy feladatok által generált táblákat is | Hivatkozzon rájuk közvetlenül, és bízza a felderítést és a hozzáférést a Unity Catalog felügyeleti és adatszármazási funkcióira. Lásd: Betöltés meglévő táblából. |
| Kis vagy statikus referenciaadatok (kereső fájlok, ritkán változó CSV-k) | Töltsd be csomagforrásként materializált nézetben. Nincs előnye annak, ha streamelsz valamit, ami alig változik. Lásd : Kis vagy statikus adathalmazok töltése felhőobjektum tárolásból. |
| Tetszőleges HTTP- vagy REST API felügyelt csatlakozó nélkül | Kérd le az adatokat az API-ból a folyamatban, vagy először mentsd le az API válaszait fájlokként. Lásd: Adatok betöltése API-ból folyamatokban. |
Minden forráshoz a következőket erősítse meg, mielőtt építesz:
- Identitás: Hogyan működik a csővezeték. A pipeline működhet szolgáltatási főként, ezért először állítsd be ezt, hogy ne támaszkodj egy személyes fiókra.
- Hálózati útvonal: Milyen kapcsolatra van szükség a forrásnak, például tároló hitelesítő engedélyre, külső helyre vagy Lakeflow Connect menedzselt kapcsolatra.
- A módosítások szemantikája: Hogyan jelzi a forrás a frissítéseket és a törléseket, ha egyáltalán jelzi ezeket. Ez határozza meg, hogy szükség van-e CDC-re, vagy a forrás csak hozzáfűzhetőként kezelhető.
Válassz fájlformátumot és tárolóréteget
A pipeline hozza meg ezt a döntést helyetted. Minden streaming tábla és a pipeline által létrehozott materializált nézet alapértelmezetten Delta táblaként van tárolva, ami ACID tranzakciókat, séma végrehajtást és fejlődést, időutazást, valamint Unity Catalog irányítást és lineage-et ad minden adathalmazon. Nem te választod ki a pipeline kimenetek formátumát. A valódi döntéseid a csővezeték két szélén vannak:
-
Nyers bemeneti formátum: Bármi, amit a forrás előállít, például CSV, JSON vagy Parquet. Az Auto Loader és a
read_files()közvetlenül támogatják ezeket. Megadd a formátumotcloudFiles.formata Python-ban vagy azformat =>argumentumot SQL-ben. Ha te irányítod a forrást, inkább a Parquetet vagy az Avro-t részesíted előnyben, mert azok viszik a sémát és jobban tömörödnek, ami felgyorsítja a felvételt és a séma következtetését. A pipeline bármelyik formátumot kezeli, ezért ne hagyd, hogy a formátum korlátozza a forrásválasztásodat. - Nyers tárolási hely: Fájlok esetén az adatokat ne egy nem felügyelt bucketútvonalra, hanem egy Unity Catalog-kötetbe helyezze, hogy az adatleszármazás és a hozzáférés-szabályozás egészen a fogadózónáig kiterjedjen. Lásd: Mik azok a Unity Catalog-kötetek?.
A pipeline által generált táblák esetében a megmaradt választási lehetőségeid a célkatalógus és séma, amelyek meghatározzák az irányítási határokat és felfedezhetőséget, valamint a nagy táblák fizikai elrendezését. Használja a CLUSTER BY (folyékony klaszterezés) funkciót, hogy a lekérdezések teljesítménye megfelelő maradjon a táblák növekedésével, a partíciók kézi finomhangolása nélkül. Lásd: Táblákhoz folyékony klaszterezés használata.
Betöltés meglévő táblából
Adatok betöltése az Azure Databricks bármely meglévő táblájából. Az adatokat átalakíthatja egy lekérdezéssel, vagy betöltheti a táblát a folyamat további feldolgozásához.
Python
@dp.table(
comment="A table summarizing counts of the top baby names for New York for 2021."
)
def top_baby_names_2021():
return (
spark.read.table("baby_names_prepared")
.filter(expr("Year_Of_Birth == 2021"))
.groupBy("First_Name")
.agg(sum("Count").alias("Total_Count"))
.sort(desc("Total_Count"))
)
SQL
CREATE OR REFRESH MATERIALIZED VIEW top_baby_names_2021
COMMENT "A table summarizing counts of the top baby names for New York for 2021."
AS SELECT
First_Name,
SUM(Count) AS Total_Count
FROM baby_names_prepared
WHERE Year_Of_Birth = 2021
GROUP BY First_Name
ORDER BY Total_Count DESC
Fájlok betöltése felhőalapú objektumtárolóból
A Databricks azt javasolja, hogy az automatikus betöltőt a folyamatokban használja a legtöbb adatbetöltési feladathoz a felhőobjektum-tárolóból vagy a Unity-katalógus kötetének fájljaiból. Az Auto Loader és a csővezetékek úgy vannak kialakítva, hogy növekményesen és idempotensen töltsék be az egyre növekvő adatokat, amint azok a felhőbeli tárolóba érkeznek. Lásd : Mi az automatikus betöltő? és adatok betöltése az objektumtárolóból.
Az alábbi példa az automatikus betöltő használatával olvassa be az adatokat a felhőbeli tárolóból.
Python
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("abfss://myContainer@myStorageAccount.dfs.core.windows.net/analysis/*/*/*.json")
)
SQL
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT *
FROM STREAM read_files(
'abfss://myContainer@myStorageAccount.dfs.core.windows.net/analysis/*/*/*.json',
format => "json"
);
Az alábbi példák az Automatikus betöltő használatával hoznak létre adathalmazokat CSV-fájlokból egy Unity Catalog-kötetben.
Python
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/Volumes/my_catalog/retail_org/customers/")
)
SQL
CREATE OR REFRESH STREAMING TABLE customers
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/retail_org/customers/",
format => "csv"
)
Megjegyzés:
- Ha az Auto Loader-t fájlértesítésekkel használja, és teljes frissítést futtat a csővezetékhez vagy a streamelési táblához, manuálisan kell eltávolítania az erőforrásokat. A CloudFilesResourceManager használatával végezhet tisztítást egy jegyzetfüzetben.
- Ha egy Unity Catalog-alapú folyamatban az Auto Loaderrel szeretne fájlokat betölteni, külső helyeket kell használnia. A Unity Catalog folyamatokkal való használatáról további információt a Unity Catalog használata folyamatokkal című témakörben talál.
Hitelesítés a felhőbeli tárolóban
Az Automatikus betöltő a Unity Catalog külső helyeit használja a felhőbeli tárolón végzett hitelesítéshez. Konfigurálnia kell egy külső helyet ahhoz a tárútvonalhoz, amelyből olvasni szeretne, és meg kell adnia a READ FILES jogosultságot a végrehajtó felhasználónak.
Az Azure Data Lake Storage-ból való betöltéshez konfiguráljon egy külső helyet, amelyet a tároló-hitelesítő adatok támogatnak, amely egy tároló konténerre hivatkozik. További információ: Csatlakozás a felhőbeli objektumtárolóhoz a Unity Catalog használatával.
Adatok betöltése üzenetbuszból
A pipeline-eket úgy konfigurálhatja, hogy az adatokat üzenetbuszokból töltse be. A Databricks azt javasolja, hogy folyamatos végrehajtással és továbbfejlesztett automatikus méretezéssel használjon streamelési táblákat, hogy a leghatékonyabb betöltést biztosítsa az üzenetbuszokról érkező alacsony késésű betöltéshez. További információkért lásd: A Lakeflow pipeline klaszterkihasználtságának optimalizálása automatikus skálázással.
Az alábbi kód például konfigurál egy streamelési táblát a Kafkából származó adatok betöltéséhez a read_kafka függvény használatával.
Python
from pyspark import pipelines as dp
@dp.table
def kafka_raw():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka_server:9092")
.option("subscribe", "topic1")
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE kafka_raw AS
SELECT *
FROM STREAM read_kafka(
bootstrapServers => 'kafka_server:9092',
subscribe => 'topic1'
);
Betöltés a Google Pub/Sub szolgáltatásból
Az alábbi példa létrehoz egy streamelési táblázatot, amely a Read_pubsub függvénnyel olvas be egy Google Pub/Sub témakörből.
Python
@dp.table
def pubsub_raw():
auth_options = {
"clientId": client_id,
"clientEmail": client_email,
"privateKey": private_key,
"privateKeyId": private_key_id
}
return (
spark.readStream
.format("pubsub")
.option("subscriptionId", "my-subscription")
.option("topicId", "my-topic")
.option("projectId", "my-project")
.options(auth_options)
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'my-subscription',
projectId => 'my-project',
topicId => 'my-topic',
clientEmail => secret('pubsub-scope', 'clientEmail'),
clientId => secret('pubsub-scope', 'clientId'),
privateKeyId => secret('pubsub-scope', 'privateKeyId'),
privateKey => secret('pubsub-scope', 'privateKey')
);
A Databricks a titkos kódok használatát javasolja az engedélyezési lehetőségek megadásakor. Tekintse meg a Pub/Sub hozzáférésének konfigurálását az összes hitelesítési beállításhoz.
Ha más üzenetbusz-forrásokból szeretne betöltést elvégezni, tekintse meg a következőt:
- Kinesis: read_kinesis
- Pulsar: read_pulsar
Adatok beolvasása az Azure Event Hubs szolgáltatásból
Az Azure Event Hubs egy adatfolyam-szolgáltatás, amely Apache Kafka-kompatibilis felületet biztosít. A folyamat futtatókörnyezetében található Strukturált streamelési Kafka-összekötő használatával betöltheti Azure Event Hubs üzeneteket. Az Azure Event Hubsból érkező üzenetek betöltéséről és feldolgozásáról további információt az Azure Event Hubs használata folyamatadatforrásként című témakörben talál.
Adatok betöltése külső rendszerekből
A folyamatok támogatják az adatok betöltését bármely, az Azure Databricks által támogatott adatforrásból. Lásd: Csatlakozás adatforrásokhoz és külső szolgáltatásokhoz. Külső adatokat is betölthet a Lakehouse Federation használatával támogatott adatforrásokhoz. Mivel a Lakehouse Federation használatához a Databricks Runtime 13.3 LTS vagy újabb verziója szükséges, a Lakehouse Federation használatához a folyamatot úgy konfigurálja, hogy az előnézeti csatornát használja.
Egyes adatforrások nem rendelkeznek egyenértékű SQL-támogatással. Ha nem tudja használni a Lakehouse Federationt ezen adatforrások egyikével, a Python használatával betöltheti az adatokat a forrásból. Python- és SQL-forrásfájlokat is hozzáadhat ugyanahhoz a folyamathoz. Az alábbi példa egy materializált nézetet deklarál egy távoli PostgreSQL-táblában lévő adatok aktuális állapotának eléréséhez.
import dp
@dp.table
def postgres_raw():
return (
spark.read
.format("postgresql")
.option("dbtable", table_name)
.option("host", database_host_url)
.option("port", 5432)
.option("database", database_name)
.option("user", username)
.option("password", password)
.load()
)
Kis méretű vagy statikus adathalmazok betöltése a felhőbeli objektumtárolóból
Kis méretű vagy statikus adathalmazokat az Apache Spark betöltési szintaxisával tölthet be. Az adatcsatornák támogatják az Azure Databricksben az Apache Spark által támogatott összes fájlformátumot. A teljes listát az Adatformátum beállításaicímű témakörben találja.
Az alábbi példák bemutatják a JSON betöltését egy tábla létrehozásához.
Python
@dp.table
def clickstream_raw():
return (spark.read.format("json").load("/databricks-datasets/wikipedia-datasets/data-001/clickstream/raw-uncompressed-json/2015_2_clickstream.json"))
SQL
CREATE OR REFRESH MATERIALIZED VIEW clickstream_raw
AS SELECT * FROM read_files(
"/databricks-datasets/wikipedia-datasets/data-001/clickstream/raw-uncompressed-json/2015_2_clickstream.json"
)
Megjegyzés:
A read_files SQL-függvény az Azure Databricks összes SQL-környezetében gyakori. Ez a folyamatokban az SQL használatával történő közvetlen fájlhozzáférés ajánlott mintája. További információ: Beállítások.
Adatok betöltése egyéni Python-adatforrásból
A Python egyéni adatforrásai lehetővé teszik az adatok egyéni formátumban való betöltését. Írhat kódot egy adott külső adatforrásból való olvasáshoz és íráshoz, vagy használhatja a meglévő Python kódot, hogy adatokat olvasson be a saját belső rendszereiből. A Python-adatforrások fejlesztésével kapcsolatos további információkért lásd: PySpark egyéni adatforrások.
Az alábbi példa egy egyéni adatforrást regisztrál a formátum nevével my_custom_datasource , és beolvassa belőle kötegelt és streamelési módban is.
from pyspark import pipelines as dp
# Assume `my_custom_datasource` is a custom Python custom data
# source that supports both batch and streaming reads, and has
# been registered using `spark.dataSource.register`.
# This creates a materialized view
@dp.table(name = "read_from_batch")
def read_from_batch():
return spark.read.format("my_custom_datasource").load()
# This creates a streaming table
@dp.table(name = "read_from_streaming")
def read_from_streaming():
return spark.readStream.format("my_custom_datasource").load()
Streamelési tábla konfigurálása a forrásstreamelési tábla módosításainak figyelmen kívül hagyásához
Alapértelmezés szerint a streamelési táblák csak hozzáfűző forrásokat igényelnek. Ha a forrásstreamelési tábla frissítéseket vagy törléseket igényel (például a GDPR "az elfelejtettséghez való jog" feldolgozásához), a jelölő használatával figyelmen kívül hagyhatja ezeket a skipChangeCommits módosításokat. Ez a jelző csak a spark.readStream függvény használatával működik option(), és nem használható, ha a forrásstreamelési tábla egy create_auto_cdc_flow() függvény célpontja. További információ: A Forrás Delta Lake-táblák módosításainak kezelése.
@dp.table
def b():
return spark.readStream.option("skipChangeCommits", "true").table("A")
Biztonságos hozzáférés a tárolási hitelesítő adatokhoz titkokkal egy csővezetékben
Az Azure Databricks titkos kulcsokat használhatja hitelesítő adatok, például hozzáférési kulcsok vagy jelszavak tárolására. A titok konfigurálásához a pipeline beállítások cluster konfigurációjában használjon egy Spark jellemzőt. Lásd a klasszikus számítási környezet konfigurálását a folyamatokhoz.
Az alábbi példa egy titkos kulcsot használ egy hozzáférési kulcs tárolásához, amely egy Azure Data Lake Storage tárfiók bemeneti adatainak automatikus betöltő használatával történő beolvasásához szükséges. Ugyanezzel a módszerrel konfigurálhatja a folyamathoz szükséges titkos kulcsokat, például az AWS-kulcsokat az S3 eléréséhez, vagy egy Apache Hive-metaadattár jelszavát.
Az Azure Data Lake Storage használatával kapcsolatos további információkért lásd: Csatlakozás az Azure Data Lake Storage-hoz és a Blob Storage-hoz.
Megjegyzés:
A titkos kulcs értékét beállító spark.hadoop. konfigurációs kulcshoz hozzá kell adnia a spark_conf előtagot.
{
"id": "43246596-a63f-11ec-b909-0242ac120002",
"storage": "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/<path>",
"clusters": [
{
"spark_conf": {
"spark.hadoop.fs.azure.account.key.<storage-account-name>.dfs.core.windows.net": "{{secrets/<scope-name>/<secret-name>}}"
},
"autoscale": {
"min_workers": 1,
"max_workers": 5,
"mode": "ENHANCED"
}
}
],
"development": true,
"continuous": false,
"libraries": [
{
"notebook": {
"path": "/Users/user@databricks.com/Pipeline Notebooks/pipeline quickstart"
}
}
],
"name": "pipeline quickstart using ADLS2"
}
Ebben a kódmintában cserélje le a következő értékeket.
| Placeholder | Csere erre: |
|---|---|
<container-name> |
A Azure tárfiók tárolójának neve. |
<storage-account-name> |
Az ADLS-tárfiók neve. |
<path> |
A folyamat kimeneti adatainak és metaadatainak elérési útja. |
<scope-name> |
Az Azure Databricks titkosítási hatókör neve. |
<secret-name> |
A Azure tárfiók hozzáférési kulcsát tartalmazó kulcs neve. |
from pyspark import pipelines as dp
json_path = "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/<path-to-input-dataset>"
@dp.create_table(
comment="Data ingested from an ADLS2 storage account."
)
def read_from_ADLS2():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load(json_path)
)
Ebben a kódmintában cserélje le a következő értékeket.
| Placeholder | Csere erre: |
|---|---|
<container-name> |
A bemeneti adatokat tároló Azure tárfiók tárolójának neve. |
<storage-account-name> |
Az ADLS-tárfiók neve. |
<path-to-input-dataset> |
A bemeneti adatkészlet elérési útja. |