Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Data z libovolného zdroje dat podporovaného Apache Sparkem v Azure Databricks můžete načíst pomocí kanálů. Datové sady (tabulky a zobrazení) v kanálu můžete definovat proti libovolnému dotazu, který vrací datový rámec Sparku, včetně streamovaných datových rámců a pandas pro datové rámce Spark. Pro úlohy příjmu dat doporučuje Databricks používat streamované tabulky pro většinu případů použití. Streamovací tabulky jsou užitečné pro příjem dat z cloudového úložiště objektů pomocí Auto Loaderu nebo ze sběrnic zpráv, jako je Kafka. Další informace o streamovaných tabulkách, primárním typem datové sady pro příjem dat, najdete v tématu Tabulky streamování.
Ne všechny zdroje dat mají podporu SQL pro příjem dat. Zdroje SQL a Python ale můžete kombinovat ve stejném kanálu, abyste mohli Python použít tam, kde je to potřeba. Podrobnosti o práci s knihovnami, které nejsou ve výchozím nastavení zabalené s kanály, najdete v tématu Správa Python závislostí pro kanály. Obecné informace o ingestování v Azure Databricks najdete v článku Zvolení standardního konektoru.
Následující příklady ukazují některé běžné vzory načítání dat.
Identifikujte své zdroje dat a cestu připojení
Než napíšete kód pipeline, inventuujte všechna místa, odkud data pocházejí. U každého zdroje si všimněte, jak zpřístupňuje svá data (soubory, databázi, systém software jako služba (SaaS), API nebo stream), jak často se mění a jaké přihlašovací údaje a přístup k síti potřebuje. Metoda připojení často určuje, zda je zdroj přirozeně dávkový nebo streamovaný, takže pokud to správně uděláte včas, vyhnete se pozdějšímu přepracování.
Každý zdroj zařaďte do jedné z následujících cest spojení. Následující tabulka uvádí preferovaný mechanismus pro každý zdroj potrubí:
| Source | Cesta připojení |
|---|---|
| Soubory přistávají v cloud object storage (S3, Azure Data Lake Storage, GCS) | Nejčastější výchozí bod. Použijte AutoLoader (cloudFiles formát), který se stará o inkrementální objevování, inferenci schémat a evoluci schémat. Viz Načtení souborů z cloudového úložiště objektů. |
| Databáze a SaaS aplikace (Salesforce, SQL Server, PostgreSQL, Workday) | Použijte spravovaný konektor Lakeflow Connect tam, kde jeden existuje, pro váš zdroj. Spravované konektory jsou řízeny konfigurací a zajišťují autentizaci a inkrementální nebo CDC extrakci za vás. Viz koncepty konektorů Lakeflow Connect. Pokud pro daný zdroj neexistuje žádný spravovaný konektor, načtěte ho přímo nebo nejprve uložte jeho odpovědi jako soubory. Vizte Příjem dat z API v kanálech. |
| Sběrnice zpráv (Kafka, Kinesis, Azure Event Hubs, Pub/Sub) | Čtěte přímo jako strukturovaný streamovací zdroj, protože jde o nativní streamovací zdroje. Viz Načíst data ze sběrnice zpráv. |
| Další tabulky Delta nebo položky katalogu Unity Catalog, včetně tabulek vytvořených jinými kanály nebo úlohami | Odkazujte na ně přímo a nechte správu a sledování původu dat v Unity Catalog zajistit vyhledávání a přístup. Viz Načítání z existující tabulky. |
| Malá nebo statická referenční data (vyhledávací soubory, zřídka měnící CSV) | Načíst jako dávkový zdroj v materializovaném zobrazení. Streamování něčeho, co se téměř nemění, nemá žádný smysl. Viz Načítání malých nebo statických datových sad z cloudového objektového úložiště. |
| Libovolné HTTP nebo REST API bez spravovaného konektoru | Nejprve stahujte z API v pipeline nebo přistávejte jeho odpovědi jako soubory. Vizte Příjem dat z API v kanálech. |
U každého zdroje si před stavbou ověřte následující:
- Identita: Jak pipeline funguje. Pipeline mohou běžet jako principal služby, proto to nastavte nejdříve, abyste se vyhnuli závislosti na osobním účtu.
- Síťová cesta: Jaké připojení zdroj potřebuje, například přihlašovací údaje k úložišti, externí lokalitu nebo spravované připojení Lakeflow Connect.
- Změna sémantiky: Jak zdroj signalizuje aktualizace a smazání, pokud vůbec. To určuje, zda potřebujete CDC, nebo můžete zdroj považovat pouze za připojený.
Vyberte formát souboru a vrstvu úložiště
Potrubí rozhoduje většinu za vás. Každá streamovaná tabulka a materializovaný pohled, který pipeline vytvoří, je ve výchozím nastavení uložen jako Delta tabulka, což vám poskytuje transakce ACID, vynucování a vývoj schématu, cestování časem a správu a rodokmen Unity Catalog na každém datovém setu. Formát pro výstupy pipeline si nevybíráte. Vaše skutečná rozhodnutí jsou na dvou okrajích procesu:
-
Surový vstupní formát: Cokoliv zdroj vyprodukuje, například CSV, JSON nebo Parquet. Auto Loader a
read_files()tyto podporují přímo. Formát zadejte v Pythonu pomocíformat =>nebo v SQL pomocí argumentucloudFiles.format. Pokud ovládáte zdroj, preferujte Parquet nebo Avro, protože lépe nesou schéma a komprimují, což urychluje příjem a inferenci schémat. Pipeline zvládají všechny tyto formáty, takže nenechte formát omezit výběr zdroje. - Umístění nezpracovaných dat: U souborů ukládejte data do svazku v Unity Catalog místo do nespravované cesty bucketu, aby se linie původu dat a řízení přístupu rozšiřovaly až do vstupní zóny. Viz Co jsou svazky katalogu Unity?.
Pro tabulky, které vytváří pipeline, zbývá zvolit cílový katalog a schéma, které určují hranice řízení a dohledatelnost, a fyzické uspořádání velkých tabulek. Použijte CLUSTER BY (liquid clustering) k udržení dobrého výkonu dotazů, když tabulky rostou, bez nutnosti ručně ladit partition. Viz Použití metody 'liquid clustering' pro tabulky.
Načtení z existující tabulky
Načtěte data z jakékoli existující tabulky v Azure Databricks. Tabulku můžete načíst pro další zpracování v datovém kanálu nebo data transformovat pomocí dotazu.
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
Načtení souborů z cloudového úložiště objektů
Databricks doporučuje používat Auto Loader v datových tocích pro většinu zátěží příjmu dat z cloudového objektového úložiště nebo ze souborů ve svazku Unity Catalog. Automatické nástroje pro nahrávání a datové kanály jsou navrženy tak, aby postupně a idempotentně načítaly stále rostoucí množství dat, jakmile dorazí do cloudového úložiště. Podívejte se na Co je automatický zavaděč? a Načtěte data z úložiště objektů.
Následující příklad čte data z cloudového úložiště pomocí Auto Loader.
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"
);
Následující příklady používají Auto Loader k vytvoření datových sad ze souborů CSV na úložišti Unity Catalog.
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"
)
Poznámka:
- Pokud používáte Auto Loader s oznámeními o souborech a spustíte úplnou aktualizaci pipeline nebo streamovací tabulky, musíte prostředky ručně vyčistit. K vyčištění můžete použít CloudFilesResourceManager v poznámkovém bloku.
- Pokud chcete načíst soubory s funkcí Auto Loader v kanálu s povoleným katalogem Unity, musíte použít externí umístění. Další informace o používání katalogu Unity s kanály najdete v tématu Použití katalogu Unity s kanály.
Ověřování v cloudovém úložišti
Auto Loader používá externí umístění katalogu Unity ke autentizaci pomocí cloudového úložiště. Musíte nakonfigurovat externí umístění pro cestu k úložišti, ze které chcete číst, a udělit READ FILES oprávnění uživateli, který úlohu spouští.
Pokud chcete integrovat z Azure Data Lake Storage, nakonfigurujte externí lokalitu, která je podložena přihlašovacími údaji úložiště a odkazuje na kontejner úložiště. Další informace najdete v tématu Připojení ke cloudovému úložišti objektů pomocí katalogu Unity.
Načtení dat ze sběrnice zpráv
Kanály můžete nakonfigurovat tak, aby ingestovali data z sběrnic zpráv. Databricks doporučuje používat tabulky streamování s průběžným spouštěním a vylepšeným automatickým škálováním, aby se zajistilo co nejefektivnější příjem dat při načítání z sběrnic zpráv s nízkou latencí. Další informace najdete v tématu Optimalizace využití clusteru kanálu Lakeflow pomocí automatického škálování.
Následující kód například nakonfiguruje streamovací tabulku pro příjem dat ze systému Kafka pomocí funkce read_kafka .
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'
);
Příjem dat z Google Pub/Sub
Následující příklad vytvoří streamovací tabulku, která čte data z tématu Google Pub/Sub pomocí funkce read_pubsub.
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')
);
Databricks doporučuje používat tajné kódy při poskytování možností autorizace. Informace o všech možnostech ověřování najdete v tématu Konfigurace přístupu k pub/Sub .
Chcete-li přijímat data z jiných zdrojů sběrnice zpráv, přečtěte si:
- Kinesis: read_kinesis
- Pulsar: read_pulsar
Načíst data z Azure Event Hubs
Azure Event Hubs je služba streamování dat, která poskytuje rozhraní kompatibilní s Apache Kafka. K načítání zpráv ze služby Azure Event Hubs můžete použít konektor Kafka pro Structured Streaming, který je součástí běhového prostředí kanálu. Další informace o načítání a zpracování zpráv ze služby Azure Event Hubs najdete v tématu Použití služby Azure Event Hubs jako zdroje dat kanálu.
Načtení dat z externích systémů
Kanály podporují načítání dat z libovolného zdroje dat podporovaného Azure Databricks. Viz Připojení ke zdrojům dat a externím službám. Externí data můžete načíst také s využitím federace Lakehouse pro podporované zdroje dat. Vzhledem k tomu, že federace Lakehouse vyžaduje Databricks Runtime 13.3 LTS nebo vyšší, pro použití Lakehouse Federation nakonfigurujte svou datovou pipeline, aby využívala kanál preview.
Některé zdroje dat nemají ekvivalentní podporu SQL. Pokud se službou Lakehouse Federation nemůžete použít některý z těchto zdrojů dat, můžete pomocí Pythonu ingestovat data ze zdroje. Zdrojové soubory Pythonu a SQL můžete přidat do stejného kanálu. Následující příklad deklaruje materializované zobrazení pro přístup k aktuálnímu stavu dat ve vzdálené tabulce PostgreSQL.
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()
)
Načtení malých nebo statických datových sad z cloudového úložiště objektů
Malé nebo statické datové sady můžete načíst pomocí syntaxe načtení Apache Sparku. Kanály podporují všechny formáty souborů podporované Apache Sparkem na Azure Databricks. Úplný seznam najdete v tématu Možnosti formátu dat.
Následující příklady ukazují načtení JSON pro vytvoření tabulky.
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"
)
Poznámka:
Funkce read_files SQL je společná pro všechna prostředí SQL v Azure Databricks. Tento vzor se doporučuje pro přímý přístup k souborům pomocí SQL v kanálech. Další informace naleznete v tématu Možnosti.
Načtení dat z vlastního zdroje dat Pythonu
Vlastní zdroje dat Pythonu umožňují načíst data ve vlastních formátech. Můžete napsat kód pro čtení a zápis do konkrétního externího zdroje dat nebo použít existující Python kód ke čtení dat z vlastních interních systémů. Další podrobnosti o vývoji zdrojů dat Pythonu najdete v tématu Vlastní zdroje dat PySpark.
Následující příklad zaregistruje vlastní zdroj dat s názvem formátu my_custom_datasource a čte z něj v dávkovém i streamovaném režimu.
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()
Konfigurace streamované tabulky tak, aby ignorovala změny ve zdrojové streamovací tabulce
Ve výchozím nastavení streamované tabulky vyžadují zdroje pouze pro přidávání. Pokud vaše zdrojová streamovací tabulka vyžaduje aktualizace nebo odstranění (například pro zpracování práva na zapomenutí GDPR), pomocí příznaku skipChangeCommits tyto změny ignorujte. Tento příznak funguje jen s spark.readStream při použití funkce option() a nelze ho použít, pokud je streamovací tabulka zdrojem pro funkci create_auto_cdc_flow(). Další informace najdete v tématu Zpracování změn zdrojových tabulek Delta Lake.
@dp.table
def b():
return spark.readStream.option("skipChangeCommits", "true").table("A")
Bezpečný přístup k přihlašovacím údajům úložiště s tajemstvími v rámci pipeline
Můžete použít Azure Databricks tajemství k ukládání přihlašovacích údajů, jako jsou přístupové klíče nebo hesla. K nastavení tajemství ve vašem přenosovém potrubí použijte vlastnost prostředí Spark v konfiguraci clusteru nastavení potrubí. Viz Konfigurace klasického výpočtu pro pipeliny.
Následující příklad používá tajný údaj k uložení přístupového klíče, který je nezbytný ke čtení vstupních dat z účtu úložiště Azure Data Lake Storage pomocí Auto Loaderu. Stejnou metodu můžete použít ke konfiguraci jakéhokoli tajného kódu vyžadovaného vaším kanálem, například klíčů AWS pro přístup k S3 nebo k heslu k metastoru Apache Hive.
Další informace o práci se službou Azure Data Lake Storage najdete v tématu Připojení ke službě Azure Data Lake Storage ablob Storage.
Poznámka:
Do konfiguračního klíče spark.hadoop., který nastaví hodnotu tajného klíče, musíte přidat předponu spark_conf.
{
"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"
}
V této ukázce kódu nahraďte následující hodnoty.
| Zástupný symbol | Nahradit za |
|---|---|
<container-name> |
Název kontejneru účtu úložiště Azure. |
<storage-account-name> |
Název účtu úložiště ADLS. |
<path> |
Cesta pro výstupní data a metadata kanálu. |
<scope-name> |
Název oboru tajného kódu Azure Databricks. |
<secret-name> |
Název klíče obsahujícího přístupový klíč účtu úložiště Azure. |
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)
)
V této ukázce kódu nahraďte následující hodnoty.
| Zástupný symbol | Nahradit za |
|---|---|
<container-name> |
Název kontejneru účtu úložiště Azure, který ukládá vstupní data. |
<storage-account-name> |
Název účtu úložiště ADLS. |
<path-to-input-dataset> |
Cesta ke vstupní datové sadě. |