Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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.
Bevilja nödvändiga molnbehörigheter på molnleverantörssidan. Kraven varierar beroende på molnleverantör. Se Konfigurera filhändelser för en extern plats.
Ange
cloudFiles.useManagedFileEventstilltruei 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_datasamlar in fält som inte matchar det aktuella schemat. Den läggs till automatiskt av Auto Loader. -
_corrupt_recordsamlar in rader som inte kan parsas alls. Aktivera den med hjälp avcolumnNameOfCorruptRecord:
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.cleanSourcefö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
deleteläge för att ta bort filer efter inmatning. - Använd
movelä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.cleanSourceom flera Auto Loader-strömmar eller andra klienter läser från samma källkatalog.- Använd
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.maxFileAgefö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.