Bästa praxis för Auto Loader

Den här sidan beskriver bästa praxis som du kan använda för att konfigurera Auto Loader så att den körs tillförlitligt, kostnadseffektivt och i stor skala för ditt användningsfall.

Dessa metodtips minskar driftkostnaderna och förhindrar vanliga problem som är svåra att diagnostisera i produktion, till exempel: onödiga LIST API-kostnader från fullständiga kataloggenomsökningar, tyst dataförlust från schemaavvikelse och omstarter av pipelinen som orsakas av felkonfiguration av kontrollpunkter.

Mer information om konfigurationsinformation för produktion finns i Konfigurera Auto Loader för produktionsarbetslaster. Mer information om övervakning och observerbarhet finns i Övervaka och observera Auto Loader.

Välj rätt körningsramverk

Det bästa körningsramverket för ditt användningsfall beror på hur mycket kontroll du behöver över pipelinen och hur mycket driftkostnader du vill hantera. För de flesta användare och produktionspipelines är Auto Loader med Lakeflow-pipelines ett bra val. Men om du behöver maximal kontroll och anpassning använder du Automatisk inläsare med strukturerad direktuppspelning. För den enklaste konfigurationen med en hanterad upplevelse använder du en hanterad LakeFlow Connector när den är tillgänglig.

Lakeflow-pipelines utökar Strukturerad direktuppspelning med automatisk skalning, datakvalitetskontroller, schemautvecklingshantering och övervakning via händelseloggen. Databricks rekommenderar Lakeflow-pipelines för de flesta produktionsinmatningsarbetsbelastningar.

Välj rätt schemaläggnings- och utlösartyp

Den bästa schemaläggnings- och utlösartypen för ditt användningsfall beror på dina svarstidskrav och mönster för filinhämtning. För de flesta användningsfall rekommenderar Databricks en utlösare för filinmatning med filhändelser aktiverade. Detta ger låg latensinmatning till låg kostnad eftersom beräkning endast körs när nya filer tas emot. De tre utlösartyperna skiljer sig åt i när och hur ofta pipelinen startar:

  • Kontinuerlig: Pipelinen körs utan avbrott. Använd endast när svarstiden under sekunden är ett hårt krav, eftersom kontinuerlig beräkning kostar mer. Koppla ihop med filhändelser.
  • Utlösare för filinkomst: Pipelinen startar när nya filer hamnar på källplatsen. Bäst för låg till medelhög svarstid eller oregelbundna filinhämtningsmönster. Kräver att filhändelser aktiveras. Se Utlösa jobb när nya filer tas emot.
  • Schemalagt: Pipelinen körs enligt ett tidsbaserat schema (till exempel varje timme). Använd när svarstidskraven är överseende (minuter till timmar). Fungerar med kataloglistning, men filhändelser minskar kostnaderna även i schemalagt läge genom att undvika fullständiga kataloggenomsökningar.

Mer information om hur du använder Trigger.AvailableNow för batchschemaläggning finns i Använda Trigger.AvailableNow och hastighetsbegränsning.

Välj rätt filidentifieringsläge

Auto Loader stöder tre lägen för filidentifiering med olika kompromisser när det gäller konfigurationskomplexitet, skalbarhet och kostnad.

Mode Konfigurationskomplexitet Scalability Cost När det bör användas
Filhändelser (rekommenderas) Låg (engångsbehörighetskonfiguration) Miljontals filer per timme Lägsta Standard för de flesta arbetsbelastningar
Klassisk filavisering Hög (21+ molnkonfigurationsalternativ) Miljontals filer per timme Medium När filhändelser inte är tillgängliga
Kataloglista None Begränsad efter katalogstorlek Högst (LIST API-kostnader) Små kataloger, engångspåfyllningar eller när säkerhetsprinciper förhindrar filhändelser

Filhändelser konsoliderar molnlagringsresurser med hjälp av en prenumeration och kö per extern plats i stället för en per dataström. Prestandaskillnaden är betydande i stor skala: kataloglistan måste genomsöka hela källkatalogen på varje utlösare, så inmatningstiden växer med katalogstorleken. Filhändelser levererar nya filmeddelanden direkt, så inmatningstiden förblir låg oavsett hur många objekt som finns i katalogen.

Aktivera filhändelser

Filhändelser kräver en engångsbehörighet för molnet och en extern plats som konfigurerats för att använda tjänsten för hanterade filhändelser. När den har konfigurerats kan alla automatiska inläsningsströmmar som läser från den externa platsen använda filhändelser utan ytterligare konfiguration.

  1. Bevilja nödvändiga molnbehörigheter på molnleverantörssidan. Kraven varierar beroende på molnleverantör. Se Konfigurera filhändelser för en extern plats.

  2. Ange cloudFiles.useManagedFileEvents till true i frågan för automatisk inläsning.

    df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.useManagedFileEvents", "true")
      .load("/path/to/data/dir"))
    

    Fullständiga installationssteg finns i Migrera till automatisk inläsare med filhändelser.

När du inte kan använda filhändelser

Du kanske inte kan använda filhändelser när:

  • Den externa platsen är inte konfigurerad med filhändelser.
  • Organisationens säkerhetsprinciper tillåter inte aktivering av filhändelser på en delad extern plats.

I dessa fall använder du klassiskt läge för filmeddelanden eller kataloglistningsläge. En fullständig jämförelse av filidentifieringslägen finns i Jämför filidentifieringslägen för automatisk inläsning.

Hantera schemautveckling

Auto Loader identifierar schemat automatiskt, men hur du konfigurerar schemautvecklingen påverkar datakomplettheten och pipeline-stabiliteten. Använd följande tabell för att välja en strategi.

Scenario Recommendation
Schemat är känt och fast Ange ett explicit schema med .schema()
Schemat är okänt, additiva ändringar förväntas schemaEvolutionMode: addNewColumns
Schemat är okänt, typändringar förväntas schemaEvolutionMode: addNewColumnsWithTypeWidening
Strikt schemakontrakt krävs schemaEvolutionMode: failOnNewColumns
Godtyckligt eller oförutsägbart schema Importera som typen Variant

När du har valt en strategi använder du följande metoder för att finjustera hur schemautvecklingen fungerar.

Använda schematips för kända fälttyper

Använd alternativet cloudFiles.schemaHints för att framtvinga typer för fält som du känner till i förväg, samtidigt som du tillåter schemainferens för andra fält.

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaHints", "id long, amount double")
  .load("/path/to/data/dir"))

Använd typbreddning för ändringar av kompatibla typer

Schemats addNewColumnsWithTypeWidening utvecklingsläge breddar automatiskt kompatibla typer (till exempel int till long) i stället för att dirigera data till _rescued_data kolumnen. Detta eliminerar behovet av jobb för efterbearbetning för att hantera enkla typomvandlingar.

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "parquet")
  .option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
  .load("/path/to/data/dir"))

Mata in som Variant typ för oförutsägbara scheman

När dina data inte överensstämmer med något specifikt schema, eller schemat ändras kontinuerligt, matar du in data som en Variant typ.

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("singleVariantColumn", "data")
  .load("/path/to/data/dir"))

Variant tillhandahåller schema-on-read vid frågetillfället men är mindre effektivt än att köra frågor mot strukturerade kolumner. Fullständig mekanik för schemainferens och utveckling finns i Konfigurera schemainferens och utveckling i Automatisk inläsning.

Hantera felaktig data- och datakvalitet

Följande metoder hjälper dig att identifiera, samla in och isolera felaktiga data innan de sprids till underordnade lager.

Aktivera _rescued_data och _corrupt_record

Auto Loader har två kolumner för att fånga upp data som inte kan tolkas korrekt vid parsning.

  • _rescued_data samlar in fält som inte matchar det aktuella schemat. Den läggs till automatiskt av Auto Loader.
  • _corrupt_record samlar in rader som inte kan parsas alls. Aktivera den med hjälp av columnNameOfCorruptRecord:
df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaHints", "_corrupt_record string")
  .option("columnNameOfCorruptRecord", "_corrupt_record")
  .load("/path/to/data/dir"))

Databricks rekommenderar columnNameOfCorruptRecord i stället för badRecordsPath för att undvika potentiella kapplöpningstillstånd som kan leda till att skadade poster missas.

Använd förväntningar i Lakeflow-pipelines för att övervaka

Ange förväntningar i Lakeflow-pipelines för att verifiera att _rescued_data och _corrupt_record är NULL under normala omständigheter. Icke-NULL-värden signalerar schemaavvikelse eller skadade data.

import dlt

@dlt.table
@dlt.expect("no rescued data", "_rescued_data IS NULL")
@dlt.expect("no corrupt records", "_corrupt_record IS NULL")
def bronze_table():
    return (spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.schemaHints", "_corrupt_record string")
        .option("columnNameOfCorruptRecord", "_corrupt_record")
        .load("/path/to/data/dir"))

Isolera skadade data

Isolera rader som innehåller data som inte kan tolkas i en separat mottagare för vidare undersökning. Detta förhindrar att skadade data sprids till underordnade lager.

import dlt

@dlt.table
def corrupt_records_sink():
    return dlt.read_stream("bronze_table").where("_corrupt_record IS NOT NULL")

@dlt.view
def clean_table():
    return dlt.read_stream("bronze_table").where("_corrupt_record IS NULL")

Annotera data med metadata från källfilen

Inkludera kolumnen _metadata i dina Auto Loader-inläsningsfrågor. Registrera åtminstone file_path och file_modification_time. På så sätt kan du spåra dataproblem tillbaka till specifika källfiler och ansluta mot cloud_files_state() för hela fillivscykeln.

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .load("/path/to/data/dir")
  .select("*", "_metadata.file_path", "_metadata.file_modification_time"))

Mer information finns i kolumnen Filmetadata.

Optimera kostnader och prestanda

Följande metoder minskar de tre huvudsakliga kostnadsdrivrutinerna för automatisk inläsning: api-anrop i molnet LIST , inaktiv beräkning och långsiktig lagringstillväxt.

  • Använda filhändelser för att minimera LIST API-kostnader: Filhändelser ger inkrementell filidentifiering, vilket eliminerar behovet av fullständiga kataloglistor vid varje körning. Det här är den enskilt mest betydelsefulla kostnadsoptimeringen för Auto Loader.

  • Använd utlösare för filinkomst för händelsedriven bearbetning: Utlösare för filinkomst startar endast din pipeline när nya filer anländer, så du betalar inte för inaktiv beräkning. Se Utlösa jobb när nya filer tas emot.

  • Arkivera bearbetade filer med cloudFiles.cleanSource: Använd cloudFiles.cleanSource för att automatiskt ta bort eller flytta bearbetade filer. Detta minskar både lagringskostnaderna och kataloglistningskostnaderna för långlivade strömmar. Fullständig information finns i Arkivera filer i källkatalogen för att sänka kostnaderna.

    • Använd delete läge för att ta bort filer efter inmatning.
    • Använd move läge för att arkivera filer till en annan plats för efterlevnad eller granskning.
    df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.cleanSource", "delete")
      .load("/path/to/data/dir"))
    

    Warning

    Aktivera inte cloudFiles.cleanSource om flera Auto Loader-strömmar eller andra klienter läser från samma källkatalog.

  • Dra nytta av prestandaförbättringar: Uppgradera till den senaste Databricks Runtime eller använd serverlös beräkning för att dra nytta av de senaste prestandaförbättringarna för automatisk inläsning.

Kontrollpunktshantering

Kontrollpunkten lagrar strömmens förlopp och filtillstånd. Att felkonfigurera eller förlora kontrollpunkten kräver en fullständig omstart, så behandla den som kritisk infrastruktur.

  • Använd aldrig livscykelprinciper för molnobjekt på kontrollpunktsplatser. Om kontrollpunktsfiler tas bort är dataströmstillståndet skadat och du måste starta om från början.
  • Använd separata kontrollpunkter för varje dataström och källkatalog.
  • Överväg cloudFiles.maxFileAge för långvariga strömmar med hög volym för att begränsa tillståndstillväxten. Använd en konservativ inställning (minst 90 dagar rekommenderas). Om du anger det här värdet för aggressivt riskerar du att bearbeta filer som Auto Loader redan har matat in om de hamnar utanför fönstret.

Fullständig information finns i Spårning av filhändelser.

Använd volymer för optimal filupptäckt med filhändelser

För bättre prestanda med filhändelser skapar du en extern volym för varje sökväg eller underkatalog som autoinläsaren läser in från. Ange volymsökvägar (till exempel /Volumes/catalog/schema/volume) till Auto Loader i stället för molnsökvägar (till exempel s3://bucket/path). Detta optimerar filidentifieringen genom ett optimerat dataåtkomstmönster.

Mer information om bästa praxis för filhändelser finns i Bästa praxis för Auto Loader med filhändelser.