Läsa in data i pipelines

Du kan läsa in data från alla datakällor som stöds av Apache Spark på Azure Databricks med hjälp av pipelines. Du kan definiera datauppsättningar (tabeller och vyer) i en pipeline utifrån valfri fråga som returnerar en Spark DataFrame, inklusive streamande DataFrames och Pandas for Spark DataFrames. För datainmatningsuppgifter rekommenderar Databricks att du använder strömningstabeller för de flesta användningsfall. Strömmande tabeller är användbara för att mata in data från molnobjektlagring med hjälp av Auto Loader eller från meddelandebussar som Kafka. Mer information om strömmande tabeller, den primära datamängdstypen för inmatning, finns i Strömmande tabeller.

Alla datakällor har inte SQL-stöd för inmatning. Du kan dock blanda SQL- och Python-källor i samma pipeline för att använda Python där det behövs. Mer information om hur du arbetar med bibliotek som inte är paketerade med pipelines som standard finns i Hantera Python beroenden för pipelines. Allmän information om datainmatning i Azure Databricks finns i Standardanslutningar i Lakeflow Connect.

I följande exempel visas några vanliga datainläsningsmönster.

Identifiera dina datakällor och anslutningsväg

Innan du skriver pipelinekod, inventera varje plats data kommer ifrån. För varje källa, notera hur den exponerar sina data (filer, en databas, ett SaaS-system, ett API eller en ström), hur ofta den ändras och vilka inloggningsuppgifter och nätverksåtkomst den behöver. Anslutningsmetoden avgör ofta om en källa är naturligt batch- eller streaming, så att få detta rätt tidigt undviker omarbetning senare.

Lägg varje källa i en av följande anslutningsvägar. Följande tabell listar den föredragna mekanismen för varje pipelinekälla:

Source Anslutningsväg
Filer som landar i molnobjektlagring (S3, Azure Data Lake Storage, GCS) Den vanligaste utgångspunkten. Använd Auto Loader (cloudFiles format), som hanterar inkrementell upptäckt, schemainferens och schemautveckling. Se Läsa in filer från molnobjektlagring.
Databaser och SaaS-applikationer (Salesforce, SQL Server, PostgreSQL, Workday) Använd en Lakeflow Connect managed connector där det finns en sådan för din källa. Managed connectors är konfigurationsstyrda och hanterar autentisering samt inkrementell extrahering eller extrahering via CDC åt dig. Se Hanterade kontakter i Lakeflow Connect. Om det inte finns någon hanterad anslutning för din källa kan du importera den direkt eller först lagra dess svar som filer. Se Inmata data från ett API i pipelines.
Meddelandebussar (Kafka, Kinesis, Azure Event Hubs, Pub/Sub) Läs direkt som en strukturerad streamingkälla eftersom dessa är inhemska streamingkällor. Se Ladda data från en meddelandebuss.
Andra Delta-tabeller eller Unity Catalog-tillgångar, inklusive tabeller producerade av andra pipelines eller jobb Referera direkt till dem och låt Unity Catalogs styrning och härstamning hantera upptäckt och åtkomst. Se Ladda från en befintlig tabell.
Små eller statiska referensdata (uppslagsfiler, sällan föränderliga CSV:er) Ladda som batchkälla i en materialiserad vy. Det finns ingen fördel med att streama något som knappt förändras. Se Ladda små eller statiska dataset från molnobjektlagring.
Ett godtyckligt HTTP- eller REST-API utan hanterad koppling Hämta från API:et i pipelinen eller lägg in svaren som filer först. Se Inmata data från ett API i pipelines.

För varje källa, kontrollera följande innan du bygger:

  • Identitet: Hur pipelinen fungerar. Pipelines kan köras med ett tjänstens huvudnamn, så konfigurera detta först för att undvika beroende av ett personligt konto.
  • Nätverksväg: Vilken anslutning källan behöver, såsom en lagringslegitimation, en extern plats eller Lakeflow Connect-hanterad anslutning.
  • Ändra semantik: Hur källan signalerar uppdateringar och borttagningar, om alls. Detta avgör om du behöver CDC eller om du kan behandla källan som endast tillägg.

Välj ett filformat och lagringslager

Pipelines fattar större delen av det här beslutet åt dig. Varje strömningstabell och materialiserad vy som en pipeline skapar lagras som standard som en Delta-tabell, vilket ger dig ACID-transaktioner, schema-kontroll och utveckling, tidsresor samt Unity Catalog-styrning och härstamning på varje dataset. Du väljer inte formatet för pipeline-utdata. Dina verkliga beslut finns i de två kanterna av pipelinen:

  • Raw-indataformat: Vad källan än producerar, såsom CSV, JSON eller Parquet. Auto Loader och read_files() stöder dessa direkt. Ange formatet i cloudFiles.format Python eller argumentet format => i SQL. Om du kontrollerar källan, välj helst Parquet eller Avro, eftersom de innehåller schemainformation och komprimeras bättre, vilket snabbar upp inläsning och schemahärledning. Pipelines hanterar något av dessa format, så låt inte formatet begränsa ditt källkodsval.
  • Rå lagringsplats: För filer, lagra data i en Unity Catalog-volym i stället för i en ostyrd sökväg i en bucket, så att datahärkomst och åtkomstkontroll även omfattar landningszonen. Se Vad är Unity Catalog-volymer?.

För de tabeller en pipeline producerar är dina återstående val målkatalogen och schemat, som sätter styrningsgränsen och upptäckbarheten, samt den fysiska layouten av stora tabeller. Använd CLUSTER BY (liquid clustering) för att hålla frågeprestandan bra när tabellerna växer utan att manuellt justera partitioner. Se Använda flytande klustring för tabeller.

Ladda från en befintlig tabell

Läs in data från en befintlig tabell i Azure Databricks. Du kan transformera data med hjälp av en fråga eller läsa in tabellen för vidare bearbetning i din pipeline.

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

Läsa in filer från molnobjektlagring

Databricks rekommenderar att du använder Auto Loader i pipelines för de flesta datainmatningsuppgifter från objektlagring i molnet eller från filer i en Unity Catalog volym. Automatisk inläsning och pipelines är utformade för att inkrementellt och idempotent läsa in ständigt växande data när de kommer till molnlagringen. Se Vad är automatisk inläsning? och Läsa in data från objektlagring.

I följande exempel läses data från molnlagring med 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"
  );

I följande exempel används Auto Loader för att skapa datamängder från CSV-filer på en Unity Catalog-volym.

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"
)

Anmärkning

  • Om du använder Auto Loader med filaviseringar och kör en fullständig uppdatering för din pipeline eller strömningstabell måste du rensa dina resurser manuellt. Du kan använda CloudFilesResourceManager i en anteckningsbok för att utföra rensning.
  • Om du vill läsa in filer med Auto Loader i en Unity Catalog-aktiverad pipeline måste du använda externa platser. Mer information om hur du använder Unity Catalog med pipelines finns i Använda Unity Catalog med pipelines.

Autentisera till molnlagring

Auto Loader använder externa lagringsplatser i Unity Catalog för att autentisera mot molnlagring. Du måste konfigurera en extern plats för den lagringssökväg som du vill läsa från och ge behörigheten READ FILES till den körbara användaren.

Om du vill mata in från Azure Data Lake Storage konfigurerar du en extern plats som backas upp av en lagringsautentiseringsuppgift som refererar till en lagringscontainer. Mer information finns i Ansluta till molnobjektlagring med Unity Catalog.

Läsa in data från en meddelandebuss

Du kan konfigurera pipelines för att mata in data från meddelandebussar. Databricks rekommenderar att du använder strömningstabeller med kontinuerlig exekvering och förbättrad autoskalning för att ge den optimala inmatningen för låg latens inläsning från meddelandebussar. Mer information finns i Optimera lakeflow-pipelineklusteranvändning med automatisk skalning.

Följande kod konfigurerar till exempel en strömmande tabell för att mata in data från Kafka med hjälp av funktionen 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'
  );

Importera från Google Pub/Sub

I följande exempel skapas en streamingtabell som läser från ett Google Pub/Sub-topic med funktionen 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 rekommenderar att du använder hemligheter när du tillhandahåller auktoriseringsalternativ. Se Konfigurera åtkomst till Pub/Sub för alla autentiseringsalternativ.

Information om hur du matar in från andra meddelandebusskällor finns i:

Ladda data från Azure Event Hubs

Azure Event Hubs är en dataströmningstjänst som tillhandahåller ett Apache Kafka-kompatibelt gränssnitt. Du kan använda Kafka-anslutningsappen för strukturerad direktuppspelning, som ingår i pipelinekörningen, för att läsa in meddelanden från Azure Event Hubs. Mer information om hur du läser in och bearbetar meddelanden från Azure Event Hubs finns i Använda Azure Event Hubs som en pipelinedatakälla.

Läsa in data från externa system

Pipelines stöder inläsning av data från alla datakällor som stöds av Azure Databricks. Se Ansluta till datakällor och externa tjänster. Du kan också ladda in externa data med Lakehouse Federation för stödda datakällor . Eftersom Lakehouse Federation kräver Databricks Runtime 13.3 LTS eller senare måste du för att använda Lakehouse Federation konfigurera din pipeline för att använda förhandsgranskningskanalen.

Vissa datakällor har inte motsvarande SQL-stöd. Om du inte kan använda Lakehouse Federation med någon av dessa datakällor kan du använda Python för att mata in data från källan. Du kan lägga till Python- och SQL-källfiler i samma pipeline. I följande exempel deklareras en materialiserad vy för att få åtkomst till det aktuella tillståndet för data i en fjärransluten PostgreSQL-tabell.

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()
  )

Läs in små eller statiska datamängder från molnobjektlagring

Du kan läsa in små eller statiska datauppsättningar med apache Spark-inläsningssyntax. Pipelines stöder alla filformat som stöds av Apache Spark på Azure Databricks. En fullständig lista finns i Alternativ för dataformat.

I följande exempel visas hur du läser in JSON för att skapa en tabell.

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"
)

Anmärkning

Funktionen read_files SQL är gemensam för alla SQL-miljöer i Azure Databricks. Det är det rekommenderade mönstret för direkt filåtkomst med SQL i pipelines. Mer information finns i Alternativ.

Läsa in data från en anpassad Python-datakälla

Med anpassade Python-datakällor kan du läsa in data i anpassade format. Du kan skriva kod för att läsa från och skriva till en specifik extern datakälla eller använda din befintliga Python kod för att läsa data från dina egna interna system. Mer information om hur du utvecklar Python-datakällor finns i PySpark-anpassade datakällor.

I följande exempel registreras en anpassad datakälla med formatnamnet my_custom_datasource och läser från den i både batch- och strömningslägen.

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()

Konfigurera en strömmande tabell för att ignorera ändringar i en källströmningstabell

Som standard kräver strömmande tabeller källor som endast tillåter tillägg. Om den strömmande källtabellen kräver uppdateringar eller raderingar (till exempel för behandling av rätten att bli bortglömd enligt GDPR), använder du flaggan skipChangeCommits för att ignorera dessa ändringar. Den här flaggan fungerar bara med spark.readStream funktionen option() och kan inte användas när källuppspelningstabellen är målet för en create_auto_cdc_flow() funktion. Mer information finns i Hantera ändringar i Delta Lake-källtabeller.

@dp.table
def b():
   return spark.readStream.option("skipChangeCommits", "true").table("A")

Kom åt lagringsuppgifter på ett säkert sätt med säkerhetsnycklar i en pipeline

Du kan använda Azure Databricks-hemligheter för att lagra autentiseringsuppgifter som åtkomstnycklar eller lösenord. Om du vill konfigurera hemligheten i din pipeline använder du en Spark-egenskap i klusterkonfigurationen för pipelineinställningar. För pipelines, se Konfigurera klassisk beräkning.

I följande exempel används en hemlighet för att lagra en åtkomstnyckel som krävs för att läsa indata från ett Azure Data Lake Storage lagringskonto med automatisk inläsning. Du kan använda samma metod för att konfigurera alla hemligheter som krävs av din pipeline, till exempel AWS-nycklar för att komma åt S3 eller lösenordet till ett Apache Hive-metaarkiv.

Mer information om hur du arbetar med Azure Data Lake Storage finns i Ansluta till Azure Data Lake Storage och Blob Storage.

Anmärkning

Du måste lägga till prefixet spark.hadoop. till spark_conf konfigurationsnyckeln som anger det hemliga värdet.

{
  "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"
}

Ersätt följande värden i det här kodexemplet.

Platshållare Ersätt med
<container-name> Namnet på containern för Azure-lagringskontot.
<storage-account-name> Namnet på ADLS-lagringskontot.
<path> Sökvägen för pipelinens utdata och metadata.
<scope-name> Namnet på Azure Databricks-hemlighetsomfånget.
<secret-name> Namnet på nyckeln som innehåller åtkomstnyckeln för Azure lagringskonto.
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)
  )

Ersätt följande värden i det här kodexemplet.

Platshållare Ersätt med
<container-name> Namnet på containern Azure lagringskonto som lagrar indata.
<storage-account-name> Namnet på ADLS-lagringskontot.
<path-to-input-dataset> Sökvägen till indatauppsättningen.