Grunderna i Spark

Grundläggande begrepp som ligger till grund för storleksändring, optimering och felsökning. Läs detta först om du är nybörjare på Spark i Fabric.

Allmänna riktlinjer och varningar

Scenario: Spark är nytt för dig. Vad är Dos och Don'ts
Användningsfall Metodtips
Använda optimerade serialiserade format Gör: Föredrar format som Avro, Parquet eller Optimized Row Columnar (ORC) eftersom de bäddar in schema, är kompakta och optimerar lagring och bearbetning. I Fabric använder du Delta-format för atomicitets-, konsistens-, isolerings- och hållbarhetsgarantier (ACID) samt prestandafördelar.
Var försiktig med XML/JSON Förlita dig inte på schemainferens för stora JSON-filer (JavaScript Object Notation) eller XML-filer (Extensible Markup Language), eftersom Spark läser hela datamängden för att härleda schema, vilket saktar ned bearbetningen och förbrukar minne intensivt.

Ange ett statiskt primärt schema när du läser JSON/XML eller använder .option("samplingRatio", 0.1) för att påskynda läsningar, men tänk på att om exemplet inte representerar den fullständiga datamängden kan läsningarna misslyckas. En säkrare metod härleder schema från ett representativt exempel och bevarar det för alla läsningar.

Undvik att parsa stora XML-filer. XML-parsning körs av naturen långsammare på grund av taggbearbetning och typkonvertering.
Optimera kopplingar och filtrering Gör: Använd kolumnrensning och filtrering på radnivå före kopplingar för att minska användningen av shuffle och minne.

Katalysatoroptimeraren hanterar automatiskt predikat-pushdown när du använder DataFrame-API:er. Undvik RDD-API:er (Resilient Distributed Dataset) eftersom de kringgår katalysatoroptimeringar.
Föredra DataFrames framför RDD:ar Gör: Använd DataFrames i stället för RDD:er för de flesta åtgärder. DataFrames använder Catalyst-optimeraren och Tungsten-körningsmotorn för effektiv körning.
Aktivera adaptiv frågekörning (AQE) Gör: Aktivera AQE för att dynamiskt optimera shuffle-partitioner och hantera skeva data automatiskt.

Minneshantering för exekutor

Scenario: Du vill förstå hanteringen av körminnet för prestandajustering.

Även om en exekverare har konfigurerats med 56 GB minne tillåter Spark inte att allt används direkt för användardata. Spark Core delar upp och hanterar körminne:

  • Reserverat minne: En fast del reserverad för system- och Spark-interna omkostnader (till exempel Java Virtual Machine (JVM), internals).

  • Användarminne: Lagrar användardefinierade funktioner (UDF), lokala variabler, datastrukturer (listor, kartor, ordlistor) och objekt som skapats under beräkningen.

  • Lagringsminne: Innehåller cachelagrade/bevarade data, sändningsvariabler och shuffle-data som kan cachelagras.

  • Körningsminne: Används för mellanliggande beräkning (skakningar, sammanslagningar, sorteringar, aggregeringar).

  • Dynamisk minnesdelning: Gränsen mellan lagrings- och körningsminnet kan flyttas. Spark kan låna minne från en region till en annan, vilket möjliggör flexibel minnesanvändning.

  • Spills: Inträffar när efterfrågan på lagrings- eller exekveringsminne överstiger tillgängligt minne efter att det lånats. Detta tvingar data till disk, vilket kan påverka prestanda.

    Diagram över Hantering och spill av Spark-minne.

OOM-fel (Out of Memory)

Scenario: Spark-jobb misslyckas med OOM-fel (Out of Memory).

Drivrutins-OOM:

OOM-fel för drivrutinen uppstår när Spark-drivrutinen överskrider sitt allokerade minne.

Vanlig orsak: drivrutinsintensiva åtgärder som collect(), countByKey()eller stora toPandas() anrop som hämtar för mycket data till drivrutinsminnet.

Åtgärd: Undvik drivrutinsintensiva åtgärder när det är möjligt. Om det inte går att undvika kan du öka drivrutinsstorleken och prestandamåttet för att hitta den optimala konfigurationen.

OOM (Executor Out of Memory):

OOM-fel för executor uppstår när en Spark-exekutor överskrider sitt allokerade minne.

Vanlig orsak: Minnes- och beräkningsintensiva omvandlingar på stora datamängder (till exempel breda kopplingar, sammansättningar, blandningar) eller cachelagrade/bevarade datauppsättningar som överskrider den körbara filens tillgängliga minne (körning + lagringsregioner).

Åtgärd: Öka körminnet om det behövs, justera Spark-minnesfraktioner (spark.memory.fraction, spark.memory.storageFraction) och spara selektivt. Se till att cachelagrade data passar in i tillgängligt minne.

Dataskjevheter

Symtom på skevhet:

  • Några uppgifter tar längre tid än andra i Spark-användargränssnittet (scenaktiviteter visar tung svans).
  • Stort mellanrum mellan median- och maxaktivitetstider i stegmått.
  • Steg med stora shuffle-läs- och skrivstorlekar för ett fåtal partitioner.

Vanliga orsaker:

  • Ojämn datadistribution för kopplings-/gruppnycklarna (snabbnycklar).
  • Felaktig partitionering eller för få partitioner för datavolymen.
  • Uppströms dataavvikelser skapar stora poster eller många null-/tomma nycklar.

Förmildrande omständighet:

  • Ompartition eller sammansling för att öka partitionsparallellitet och balansstorlekar.
  • Använd nyckelsaltning eller anpassad partitionering för att fördela heta nycklar mellan partitioner.
  • Använd AQE (Adaptive Query Execution) för att sammanfoga partitioner efter blandning och aktivera optimering av skev koppling.
  • Använd sändningskopplingar för små uppslagstabeller för att undvika omfördelningar helt.
  • Spara balanserade mellanliggande datamängder före kostsamma processer och kör om jobbet.

Bästa praxis för UDF

Scenario: Du måste använda anpassad logik som inte kan uttryckas via inbyggda DataFrame-funktioner.

Använd Spark DataFrame-API:er när det är möjligt. Catalyst-optimeraren optimerar inbyggda funktioner och kör dem internt på JVM, så att de ger bästa möjliga prestanda.

Om du måste använda en UDF (användardefinierad funktion) bör du undvika vanliga PySpark Python-UDF:er. Tänk i stället på följande alternativ:

  • Pandas UDF:er (även kallade vektoriserade UDF:er): Använd Apache Arrow för effektiv dataöverföring mellan JVM och Python. Pandas UDF:er tillåter vektoriserade åtgärder, vilket avsevärt förbättrar prestanda jämfört med python-UDF:er rad för rad.

  • Scala/Java UDF:er: Kör direkt på JVM och undvik python-serialiseringskostnader. Scala/Java UDF:er överträffar vanligtvis Python-UDF:er.

Var försiktig med Python-UDF:er. Varje köre startar en separat Python-process som kräver serialisering och deserialisering av data mellan JVM och Python. Detta skapar en flaskhals för prestanda, särskilt i stor skala. 

Felloggning

Scenario: Bästa praxis för felloggning i Fabric Spark
  1. Använd log4j i stället för print() som belastar föraren kraftigt. Med log4jkan du komma åt loggar i drivrutinsloggar och söka efter dem (med hjälp av loggningsnamnet, till exempel: PySparkLogger).

    Diagram över Spark-loggar.

  2. Omslut läsningar, skrivningar och transformationer i try- och except-block. Använd logger.error för undantag och logger.info för förloppsmeddelanden.

    • Python-loggning: Perfekt för loggningsåtgärder, statusuppdateringar eller felsökning av information från kod som endast körs på Spark-drivrutinen. Pythons loggningsmodul sprids inte till körloggar. Se dokumentationen för att utveckla, köra och hantera notebook-filer.

    • Spark log4j: Standarden för robust programloggning på produktionsnivå i Spark eftersom den integreras internt med Sparks drivrutins-/körloggar.

    Exempel på log4j-användning i PySpark:

    import traceback
    # Get log4j logger
    log4jLogger = spark._jvm.org.apache.log4j
    logger = log4jLogger.LogManager.getLogger("PySparkLogger")
    logger.info("Application started.")
    try:
        # Create DataFrame with 20 records
        data = [(f"Name{i}", i) for i in range(1, 21)]  # 20 records
        df = spark.createDataFrame(data, ["name", "age"])
        logger.info("DataFrame created successfully with 20 records.")
        df.show(s)  # 's' is not defined -> will throw error but the application will not fail
    except Exception as e:
        logger.error(f"Error while creating or showing DataFrame: {str(e)}\n{traceback.format_exc()}")
    
  3. Centralisera felövervakning:

    • Använd diagnostikutfärdartillägget (Övervaka Apache Spark-applikationer med Azure Log Analytics) i miljön och anslut till Notebooks som kör Spark-applikationer. Emittern kan skicka händelseloggar, anpassade loggar (till exempel log4j) och mått till Azure Log Analytics/Azure Storage/Azure Event Hubs. Skicka log4j-namnet till egenskapen: spark.synapse.diagnostic.emitter.\<destination\>.filter.loggerName.match.

    • För felsökning kan du också samla in misslyckade rader/poster till Lakehouse-tabeller (LH) för felaktig datainsamling på postnivå.