Sparkin perusteet

Mitoituksen, optimoinnin ja vianmäärityksen perustana olevat ydinkäsitteet. Lue tämä ensin, jos olet uusi Spark in Fabric -käyttäjä.

Yleisiä ohjeita ja kieltoja

Skenaario: Olet uusi Spark-käyttäjä. Mitä saa ja ei saa tehdä?
Käyttötapaus Parhaat käytännöt
Käytä optimoituja sarjoitettuja muotoja Tee näin: Suosi muotoja, kuten Avro, Parquet tai Optimized Row Columnar (ORC), koska ne upottavat rakenteen, ovat kompakteja ja optimoivat tallennuksen ja käsittelyn. Käytä kankaassa Delta-muotoa atomisuuden, johdonmukaisuuden, eristyksen, kestävyyden (ACID) takuiden ja suorituskykyetujen saavuttamiseksi
Ole varovainen XML/JSON:n kanssa Älä luota skeeman päättelyyn suurissa JavaScript Object Notation (JSON) -tiedostoissa tai XML (Extensible Markup Language) -tiedostoissa, koska Spark lukee koko tietojoukon päätelläkseen skeeman, mikä hidastaa käsittelyä ja kuluttaa muistia intensiivisesti.

Anna staattinen ensisijainen rakenne, kun luet JSON/XML:ää, tai käytä .option("samplingRatio", 0.1) lukujen nopeuttamiseen, mutta muista, että jos malli ei edusta koko tietojoukkoa, lukeminen saattaa epäonnistua. Turvallisempi lähestymistapa päättelee rakenteen edustavasta otoksesta ja säilyttää sen kaikilla lukemisilla.

Vältä suurten XML-tiedostojen jäsentämistä. XML-jäsennys toimii luonnostaan hitaammin nimiöiden käsittelyn ja tyyppien muuntelun vuoksi.
Optimoi liitokset ja suodatus Tee näin: Käytä sarakkeen karsimista ja rivitason suodatusta ennen liitoksia vähentääksesi sekoitusta ja muistin käyttöä.

Catalyst-optimointi käsittelee automaattisesti predikaatin pushdownin, kun käytät DataFrame-ohjelmointirajapintoja. Vältä Resilient Distributed Dataset (RDD) -ohjelmointirajapintoja, koska ne ohittavat Catalyst-optimoinnit.
DataFrame-kehysten suosiminen RDD:iden sijaan Tee näin: Käytä tietokehyksiä RDD:iden sijaan useimmissa toiminnoissa. DataFrame-kehykset käyttävät Catalyst-optimoijaa ja Tungsten-suoritusmoottoria tehokkaaseen suoritukseen.
Ota mukautuvan kyselyn suorittaminen (AQE) käyttöön Tee näin: Ota AQE käyttöön, jotta voit dynaamisesti optimoida satunnaiset osiot ja käsitellä vääristyneitä tietoja automaattisesti.

Suorittimen muistin hallinta

Skenaario: Haluat ymmärtää suorittimen muistin hallinnan suorituskyvyn säätämistä varten.

Vaikka suorittajalle olisi määritetty 56 Gt muistia, Spark ei salli kaiken käyttöä suoraan käyttäjätietoihin. Spark Core jakaa ja hallitsee suorittimen muistia:

  • Varattu muisti: Kiinteä osa, joka on varattu järjestelmän ja Sparkin sisäisille yleiskustannuksille (esimerkiksi Java-virtuaalikone (JVM), sisäiset).

  • Käyttäjän muisti: Tallentaa käyttäjän määrittämät funktiot (UDF), paikalliset muuttujat, tietorakenteet (luettelot, kartat, sanakirjat) ja laskennan aikana luodut objektit.

  • Tallennus Muisti: Sisältää välimuistiin tallennetut/säilytetyt tiedot, lähetysmuuttujat ja välimuistiin tallennettavat sekoitustiedot.

  • Suoritusmuisti: Käytetään välilaskennassa (sekoitukset, liitokset, lajittelut, koosteet).

  • Dynaaminen muistin jakaminen: Tallennus- ja suoritusmuistin välinen raja on siirrettävissä. Spark voi lainata muistia alueelta toiselle, mikä mahdollistaa joustavan muistin käytön.

  • Läikkyä: Tapahtuu, kun joko tallennustilan tai suorituksen muistin tarve ylittää lainaamisen jälkeen käytettävissä olevan muistin. Tämä pakottaa tiedot levylle, mikä voi vaikuttaa suorituskykyyn.

    Kaavio Spark-muistin hallinnasta ja vuotamisesta.

Muistin loppuminen (OOM) -virheet

Skenaario: Spark-työt epäonnistuvat OOM-virheiden vuoksi.

Ohjaimen OOM:

Ohjaimen OOM-virheitä ilmenee, kun Spark-ohjain ylittää varatun muistin.

Yleinen syy: ohjainraskaat toiminnot, kuten collect(), countByKey()tai suuret toPandas() kutsut, jotka vetävät liikaa tietoja ohjaimen muistiin.

Lievennys: Vältä kuljettajan raskaita toimintoja aina kun mahdollista. Jos se on väistämätöntä, suurenna ohjaimen kokoa ja vertailuarvoa optimaalisen kokoonpanon löytämiseksi.

Suorittajan muisti loppunut (OOM):

Suorittimen OOM-virheitä ilmenee, kun Spark-suorittaja ylittää varatun muistin.

Yleinen syy: Muisti- ja laskentaintensiiviset muunnokset suurissa tietojoukoissa (esimerkiksi laajat liitokset, koosteet, sekoitukset) tai välimuistissa oletetuissa/säilytetyissä tietojoukoissa, jotka ylittävät suorittajan käytettävissä olevan muistin (suoritus + tallennusalueet).

Lievennys: Lisää suorittimen muistia tarvittaessa, säädä Spark-muistin murto-osat (spark.memory.fraction, spark.memory.storageFraction) ja säilytä valikoivasti. Varmista, että välimuistiin tallennetut tiedot mahtuvat käytettävissä olevaan muistiin.

Tiedot vääristyvät

Vinouman oireet:

  • Jotkin tehtävät kestävät kauemmin kuin toiset Spark-käyttöliittymässä (vaihetehtävissä on raskas häntä).
  • Suuri ero mediaani- ja maksimitehtäväaikojen välillä vaihemittareissa.
  • Vaiheet, joissa on suuret satunnaisluku- tai kirjoituskoot muutamalle osiolle.

Yleiset syyt:

  • Liittymis-/ryhmäavainten (pikanäppäimien) tietojen epätasainen jakautuminen.
  • Virheellinen osiointi tai liian vähän osioita datataltiolle.
  • Ylävirran tietojen poikkeamat, jotka tuottavat suuria tietueita tai useita tyhjä- tai tyhjiä avaimia.

Lieventäminen:

  • Osioi uudelleen tai yhdistä osion yhdensuuntaisuuden ja tasapainon koon lisäämiseksi.
  • Käytä avainten suolausta tai mukautettua osiointia levittääksesi pikanäppäimiä osioiden välillä.
  • Käytä AQE:tä (Adaptive Query Execution) sekoituksen jälkeisten osioiden yhdistämiseen ja vinoliitoksen optimointien käyttöönottoon.
  • Käytä lähetysliitoksia pienissä hakutaulukoissa, jotta vältät sekoitukset kokonaan.
  • Säilytä tasapainotetut välitietojoukot ennen kalliita vaiheita ja suorita työ uudelleen.

UDF:n parhaat käytännöt

Skenaario: Sinun on käytettävä mukautettua logiikkaa, jota ei voi ilmaista sisäisten DataFrame-funktioiden avulla.

Käytä Spark DataFrame -ohjelmointirajapintoja aina kun mahdollista. Catalyst-optimointi optimoi sisäänrakennetut toiminnot ja suorittaa ne natiivisti JVM:ssä, jotta ne tarjoavat parhaan suorituskyvyn.

Jos sinun on käytettävä UDF:ää (User Defined Function), vältä tavallisia PySpark Python UDF:iä. Harkitse sen sijaan seuraavia vaihtoehtoja:

  • Pandas UDF:t (tunnetaan myös nimellä vektorisoidut UDF:t): Käytä Apache Arrow'ta tehokkaaseen tiedonsiirtoon JVM:n ja Pythonin välillä. Pandas UDF:t mahdollistavat vektorisoidut toiminnot, mikä parantaa suorituskykyä merkittävästi rivi riviltä Python UDF:iin verrattuna.

  • Scala/Java-UDF:t: Suorita suoraan JVM:ssä, jolloin vältetään Pythonin sarjoittamisen kuormitus. Scala/Java UDF:t ovat tyypillisesti parempia kuin Python UDF:t.

Ole varovainen Python UDF:ien kanssa. Jokainen suorittaja käynnistää erillisen Python-prosessin, joka edellyttää tietojen sarjoittamista ja deserialisointia JVM:n ja Pythonin välillä. Tämä luo suorituskyvyn pullonkaulan erityisesti mittakaavassa. 

Virheiden kirjaaminen

Skenaario: Fabric Sparkin virheiden kirjaamisen parhaat käytännöt
  1. Käytä log4j sen sijaan print() kuormittaa kuljettajaa raskaasti. Painikkeella log4jvoit käyttää ohjainlokien lokeja ja hakea niitä (käyttämällä loggerin nimeä, esimerkiksi: PySparkLogger).

    Kaavio Spark-lokeista.

  2. Kääri luku-, kirjoitus- ja muunnokset try- ja except -lohkoihin. logger.error Käytä poikkeuksiin ja logger.info edistymisviesteihin.

    • Pythonin kirjaaminen: Ihanteellinen toimintojen, tilapäivitysten kirjaamiseen tai virheenkorjaukseen koodista, joka suoritetaan vain Spark-ohjaimessa. Pythonin kirjausmoduuli ei leviä suorittajalokeihin. Katso muistikirjojen kehittämisen, suorittamisen ja hallinnan dokumentaatio.

    • Kipinä log4j: Standardi vankalle tuotantotason sovellusten kirjaamiselle, koska se integroituu natiivisti Sparkin ohjain-/suorittajalokeihin.

    Esimerkki log4j:n käytöstä PySparkissa:

    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. Keskitä virheiden seuranta:

    • Käytä diagnostiikkalähetinlaajennusta (Apache Spark -sovellusten valvonta Azure Log Analyticsin avulla) ympäristössä ja liitä Spark-sovelluksia suorittaviin muistikirjoihin. Lähettäjä voi lähettää tapahtumalokeja, mukautettuja lokeja (kuten log4j) ja mittareita Azure Log Analyticsiin/Azure Storageen/Azuren tapahtumatoimintoihin. Välitä log4j-nimi kiinteistölle: spark.synapse.diagnostic.emitter.\<destination\>.filter.loggerName.match.

    • Lisäksi virheenkorjausta varten voit myös kerätä epäonnistuneita rivejä/tietueita Lakehouse (LH) -taulukoihin tietuetason virheellisten tietojen keräämistä varten.