Parse 25TB me AWK dhe R

Parse 25TB me AWK dhe R
Si e lexoni kĂ«tĂ« artikull: mĂ« falni qĂ« teksti doli kaq i gjatĂ« dhe kaotik. PĂ«r tĂ« kursyer kohĂ«n tuaj, çdo kapitull e filloj me njĂ« hyrje "ÇfarĂ« mĂ«sova", ku pĂ«rmbledh nĂ« njĂ« ose dy fjali thelbin e kapitullit.

"Thjesht trego zgjidhjen!" Nëse dëshironi vetëm të shihni se çfarë arrita, kaloni në kapitullin "Po bëhem më shpikës", por unë mendoj se është më interesante dhe e dobishme të lexoni për dështimet.

Së fundmi, më kërkuan të konfiguroj procesin e përpunimit të një volum të madh të sekuencave të ADN-së (teknikisht ky është SNP-chip). Duhej të merrja shpejt të dhëna për një lokacion të caktuar gjenetik (i cili quhet SNP) për modelimin e mëtejshëm dhe detyra të tjera. Me ndihmën e R dhe AWK, arrita të pastroj dhe organizoj të dhënat në një mënyrë natyrale, duke e përshpejtuar ndjeshëm përpunimin e kërkesave. Kjo nuk ishte e lehtë për mua dhe kërkoi shumë iteracione. Ky artikull do t'ju ndihmojë të shmangni disa nga gabimet e mia dhe do të demonstrojë se çfarë arrita në fund.

Së pari, disa shpjegime hyrëse.

Të dhënat

Qendra jonĂ« universitare pĂ«r pĂ«rpunimin e informacionit gjenetik na ofroi tĂ« dhĂ«nat nĂ« formatin TSV me volum 25 TB. Ato mĂ« erdhĂ«n tĂ« ndara nĂ« 5 paketa, tĂ« kompresuara me Gzip, secila prej tĂ« cilave pĂ«rmbante rreth 240 skedarĂ« katĂ«r gigabajtĂ«sh. Çdo rresht pĂ«rmbante tĂ« dhĂ«na pĂ«r njĂ« SNP tĂ« njĂ« personi. NĂ« total u dorĂ«zuan tĂ« dhĂ«na pĂ«r ~2.5 milion SNP dhe ~60 mijĂ« njerĂ«z. PĂ«rveç informacionit SNP, nĂ« skedarĂ« kishte shumĂ« kolona me numra qĂ« pasqyronin karakteristika tĂ« ndryshme, si intensiteti i leximit, frekuenca e aleleve tĂ« ndryshme etj. NĂ« total kishte rreth 30 kolona me vlera unike.

Qëllimi

Ashtu si në çdo projekt menaxhimi të të dhënave, më e rëndësishmja ishte të përcaktohej se si do të përdoren të dhënat. Në këtë rast ne kryesisht do të përshtatim modele dhe flukse pune për SNP në bazë të SNP. Pra, në të njëjtën kohë, na nevojiten të dhëna vetëm për një SNP. Duhej të mësoja si të nxjerr të gjitha shënimet që i përkasin njërit prej 2.5 milion SNP-ve sa më thjeshtë, shpejt dhe lirë.

Si të mos e bëni këtë

Do të citoj një klishe të përshtatshme:

Nuk kam dështuar mijëra herë, thjesht kam zbuluar mijëra mënyra për të mos zhvilluar një sasi të madhe të dhënash në një format të përshtatshëm për kërkesa.

Përpjekja e parë

ÇfarĂ« mĂ«sova: nuk ka njĂ« mĂ«nyrĂ« tĂ« lirĂ« pĂ«r tĂ« analizuar 25 Tb njĂ«herĂ«sh.

Pas dëgjimit të kursit "Metodat Avancuara të Përpunimit të Dhënave të Mëdha" në Universitetin Vanderbilt, isha i sigurt se isha në rrugën e duhur. Ndoshta do të duheshin rreth një orë deri dy për të konfiguruar një server Hive për të kaluar nëpër të dhëna dhe për të raportuar rezultatin. Duke qenë se të dhënat tona ruhen në AWS S3, përfitoi nga shërbimi Athena, i cili lejon aplikimin e kërkesave Hive SQL në të dhënat S3. Nuk është e nevojshme të konfigurosh/nxjerrësh një klaster Hive, dhe paguani vetëm për të dhënat që kërkoni.

Pasi i tregova Athena të dhënat e mia dhe formatin e tyre, kam ekzekutuar disa teste me kërkesa të ngjashme:

select * from intensityData limit 10;

Dhe shpejt mora rezultate të strukturuara mirë. Gati.

Derisa nuk e provuam të përdorim të dhënat në punë...

MĂ« kĂ«rkuan tĂ« nxjerr gjithĂ« informacionin mbi SNP, pĂ«r ta testuar atĂ« nĂ« model. ÇfarĂ« kĂ«rkoi:


select * from intensityData 
where snp = 'rs123456';

...dhe fillova tĂ« prisja. Pasi kaluan tetĂ« minuta dhe mĂ« shumĂ« se 4 Tb tĂ« dhĂ«nash tĂ« kĂ«rkuara, mora rezultatin. Athena ka njĂ« tarifĂ« pĂ«r volumin e tĂ« dhĂ«nave tĂ« gjetura, $5 pĂ«r terabajt. Pra, ky kĂ«rkesĂ« i vetĂ«m mĂ« kushtoi $20 dhe tetĂ« minuta pritje. PĂ«r tĂ« kaluar modelin pĂ«r tĂ« gjitha tĂ« dhĂ«nat, do tĂ« duhej tĂ« prisja 38 vjet dhe tĂ« paguaj $50 milion. ËshtĂ« e qartĂ« se kjo nuk ishte zgjidhja pĂ«r ne.

Duhej të përdorim Parquet...

ÇfarĂ« mĂ«sova: kujdes nga madhĂ«sia e skedarĂ«ve tuaj Parquet dhe organizimi i tyre.

Fillimisht përpiqesha të rregulloja situatën duke konvertuar të gjitha TSV në skedarë Parquet.Ato janë të përshtatshme për punë me grupe të mëdha të dhënash, sepse informatat në to ruhen në format kolonor: çdo kolonë ndodhet në një segment të vetëm të memories/diskut, ndryshe nga skedarët tekstualë, ku rreshtat përmbajnë elemente të çdo kolone. Dhe nëse duhen gjetur disa, mjafton të lexoni kolonën përkatëse. Për më tepër, në çdo skedar në kolonë ruhet një gamë vlerash, prandaj, nëse vlera e kërkuar nuk është në gamën e kolonës, Spark nuk do të humbasë kohë duke skanuar të gjithë skedarin.

Kam nisur një detyrë të thjeshtë AWS Glue për të konvertuar TSV-të tona në Parquet dhe hedhur skedarë të rinj në Athena. Kjo zbuloi rreth 5 orë. Por kur e nisi kërkesën, ekzekutimi i saj zgjati afërsisht po aq dhe disi më pak para. Problemi është se Spark, duke përpjekur të optimizojë detyrën, thjesht e çoi një TSV-chunk dhe e vendosi në një chunk të tijin Parquet. Dhe, pasi çdo chunk ishte mjaft i madh dhe përmbante të dhëna të plota për shumë njerëz, çdo skedar kishte të gjitha SNP-të, kështu që Spark duhej të hapte të gjitha skedarët për të marrë informacionin e nevojshëm.

ËshtĂ« e çuditshme qĂ« lloji i kompresimit i pĂ«rdorur nga default (dhe i rekomanduar) nĂ« Parquet — snappy, — nuk Ă«shtĂ« ndarĂ«s (splitable). Prandaj çdo execucator ngec nĂ« detyrĂ«n e ç'kompresimit dhe ngarkimi i dataset-it tĂ« plotĂ« prej 3.5 GB.

Parse 25TB me AWK dhe R

Po shqyrtojmë problemin

ÇfarĂ« mĂ«sova: Ă«shtĂ« e komplikuar tĂ« renditĂ«sh, veçanĂ«risht nĂ«se tĂ« dhĂ«nat janĂ« tĂ« shpĂ«rndara.

Më dukej se tani kuptoja thelbin e problemit. Ndoshta duhej thjesht të renditja të dhënat sipas kolonës SNP, jo sipas njerëzve. Në këtë mënyrë, në një chunk të veçantë të të dhënave do të ishin disa SNP, dhe atëherë do të shfaqej në të gjithë madhështinë e saj funksioni "inteligjent" të Parquet-it "hap vetëm nëse vlera është brenda gamës". Fatkeqësisht, renditja e miliardave rreshtash të shpërndara në një grup rezultoi të ishte një detyrë sfiduese.

AWS me siguri nuk dëshiron të kthejë paratë për shkak të "Jam një student i shpërndarë". Pasi e nisa renditjen në Amazon Glue, ajo punoi për 2 ditë dhe përfundoi me dështim.

ÇfarĂ« ndodhi me ndarjen nĂ« pjesĂ«?

ÇfarĂ« mĂ«sova: ndarjet nĂ« Spark duhet tĂ« jenĂ« tĂ« balancuara.

Pastaj më erdhi ideja për të ndarë të dhënat në kromozome. Janë 23 (dhe disa të tjera nëse merren parasysh DNA mitokondriale dhe territore të paqartësuara (unmapped)).
Kjo do të lejojë ndarjen e të dhënave në pjesë më të vogla. Nëse shtojmë në funksionin e eksportit të Spark në skriptin Glue vetëm një rresht partition_by = "chr", atëherë të dhënat duhet të jenë të ndara në bakete (buckets).

Parse 25TB me AWK dhe R
Gjenomi përbëhet nga shumë fragmente, të quajtur kromozome.

Fatkeq, kjo nuk funksionoi. Kromozomet kanë madhësi të ndryshme, kështu që kanë edhe sasi të ndryshme informacioni. Kjo do të thotë se detyrat që Spark i dërgoi punëtorëve nuk ishin të balancuara dhe u ekzekutuan ngadalë, pasi disa nyje përfunduan më herët dhe mbetën pa punë. Megjithatë, detyrat u përfunduan. Por kur kërkohej një SNP, disbalancimi përsëri shkaktoi probleme. Kostumi i përpunimit të SNP në kromozomet më të mëdha (pra atje ku ne duam të marrim të dhënat) u ul vetëm rreth 10 herë. Shumë, por jo mjaftueshëm.

Dhe çfarë nëse e ndajmë në parti edhe më të vogla?

ÇfarĂ« mĂ«sova: asnjĂ«herĂ« mos u pĂ«rpiqni tĂ« bĂ«ni 2.5 milion partie.

UnĂ« vendosa tĂ« ecja deri nĂ« fund dhe i partiçionova çdo SNP. Kjo garantoi madhĂ«si tĂ« njĂ«jtĂ« pĂ«r partitĂ«. ISHTE NJË IDE E KEQE. PĂ«rdora Glue dhe shtova njĂ« varg tĂ« pafajshĂ«m partition_by = 'snp'. Detyra u nis dhe filloi tĂ« ekzekutohej. NjĂ« ditĂ« mĂ« vonĂ« kontrollova dhe pashĂ« se nĂ« S3 ende nuk ishte shkruar asgjĂ«, prandaj e ndalova detyrĂ«n. Duket se Glue shkruante skedarĂ«t pĂ«rkohĂ«sorĂ« nĂ« njĂ« vend tĂ« fshehur nĂ« S3, dhe kishte shumĂ« skedarĂ«, ndoshta disa milion. Si rezultat, gabimi im kushtoi mĂ« shumĂ« se njĂ« mijĂ« dollarĂ« dhe nuk e gĂ«zoi mentorin tim.

Particionimi + renditja

ÇfarĂ« mĂ«sova: ende Ă«shtĂ« e vĂ«shtirĂ« tĂ« renditĂ«sh, ashtu si dhe tĂ« konfigurosh Spark.

PĂ«rpjekja e fundit pĂ«r partitionim ishte qĂ« unĂ« partitionova kromozomet dhe mĂ« pas rendita çdo parti. NĂ« teori, kjo do tĂ« lejonte tĂ« pĂ«rshpejtonte çdo kĂ«rkesĂ«, sepse tĂ« dhĂ«nat e dĂ«shiruara pĂ«r SNP do tĂ« ishin brenda disa Parquet chank-Ă«ve brenda njĂ« diapazoni tĂ« dhĂ«nĂ«. FatkeqĂ«si, renditja e tĂ« dhĂ«nave tĂ« partiçionuara rezultoi tĂ« ishte njĂ« detyrĂ« e vĂ«shtirĂ«. Si rezultat, kalova nĂ« EMR pĂ«r njĂ« grup tĂ« personalizuar dhe pĂ«rdora tetĂ« instanca tĂ« fuqishme (C5.4xl) dhe Sparklyr pĂ«r tĂ« krijuar njĂ« proces mĂ« fleksibĂ«l


# Sparklyr snippet to partition by chr and sort w/in partition
# Join the raw data with the snp bins
raw_data
  group_by(chr) %>%
  arrange(Position) %>% 
  Spark_write_Parquet(
    path = DUMP_LOC,
    mode = 'overwrite',
    partition_by = c('chr')
  )


sidoqoftĂ« detyra pĂ«rsĂ«ri nuk u ekzekutua. UnĂ« e konfigurova nĂ« tĂ« gjitha mĂ«nyrat: rita memorie pĂ«r çdo ekzekutor tĂ« kĂ«rkesave, pĂ«rdora nyje me mĂ« shumĂ« memorie, aplikoja variabla tĂ« transmetimit (broadcasting variable), por çdo herĂ«, kĂ«to ishin gjysmĂ« masa, dhe gradualisht ekzekutorĂ«t filluan tĂ« prisnin, derisa gjithçka ndaloi.

Po bëhem më kreativ

ÇfarĂ« mĂ«sova: ndonjĂ«herĂ« tĂ« dhĂ«nat e veçanta kĂ«rkojnĂ« zgjidhje tĂ« veçanta.

Çdo SNP ka njĂ« vlerĂ« pozite. Ky numĂ«r korrespondon me numrin e bazave qĂ« ndodhen pĂ«rgjatĂ« kromozomit tĂ« tij. Kjo Ă«shtĂ« njĂ« mĂ«nyrĂ« e mirĂ« dhe natyrale pĂ«r tĂ« organizuar tĂ« dhĂ«nat tona. Fillimisht doja tĂ« bĂ«ja ndarje sipas rajoneve tĂ« çdo kromozomi. PĂ«r shembull, pozitat 1 - 2000, 2001 - 4000, etj. Por problemi Ă«shtĂ« se SNP janĂ« tĂ« shpĂ«rndara nĂ« kromozome nĂ« mĂ«nyrĂ« tĂ« pabarabartĂ«, kĂ«shtu qĂ« madhĂ«sia e grupeve do tĂ« ndryshonte ndjeshĂ«m.

Parse 25TB me AWK dhe R

Si rezultat, arrij në ndarjen sipas kategorive (rank) të pozita. Nga të dhënat e ngarkuara tashmë, ekzekutova një pyetje për të marrë një listë të SNP-ve unikë, pozitat dhe kromozomet e tyre. Pastaj rendita të dhënat brenda çdo kromozomi dhe grumbullova SNP në grupe (bin) me madhësi të caktuar. Le të themi, për 1000 SNP. Kjo më dha lidhjen e SNP me grupin-në-kromozom.

Në fund, bëra grupe (bin) me 75 SNP, arsyen e shpjegoj më poshtë.

snp_to_bin % 
  group_by(chr) %>% 
  arrange(position) %>% 
  mutate(
    rank = 1:n()
    bin = floor(rank/snps_per_bin)
  ) %>% 
  ungroup()

Përpjekja e parë me Spark

ÇfarĂ« mĂ«sova: bashkimi nĂ« Spark funksionon shpejt, por ndarja ende kushton shumĂ«.

Doja të lexoja këtë kuadër të vogël (2.5 milion rreshta) në Spark, ta bashkoja me të dhëna të papërpunuara, dhe pastaj ta ndaja sipas kolonës së sapo shtuar. bin.


# Join the raw data with the snp bins
data_w_bin <- raw_data %>%
  left_join(sdf_broadcast(snp_to_bin), by ='snp_name') %>%
  group_by(chr_bin) %>%
  arrange(Position) %>% 
  Spark_write_Parquet(
    path = DUMP_LOC,
    mode = 'overwrite',
    partition_by = c('chr_bin')
  )

Kam përdorur sdf_broadcast(), kështu që Spark merr vesh se duhet të dërgojë kuadrin e të dhënave në të gjitha nyjet. Kjo është e dobishme nëse të dhënat janë të vogla dhe kërkohen për të gjitha detyrat. Ndryshe, Spark përpiqet të jetë i mençur dhe shpërndan të dhënat sipas nevojës, që mund të shkaktojë ngadalësim.

Dhe përsëri, ideja ime nuk funksionoi: detyrat punuan për njëfarë kohe, përfunduan bashkimin, dhe pastaj, siç bënë ekzekutuesit që ishin nën ndarjen, filluan të dështojnë.

Po shtoj AWK

ÇfarĂ« mĂ«sova: mos flini kur ju mĂ«sojnĂ« bazat. Sigurisht, dikush e ka zgjidhur problemin tuaj ndonjĂ«herĂ« nĂ« vitet 1980.

Derisa të këtij momenti, arsyeja e të gjitha dështimeve të mia me Spark ishte përzierja e të dhënave në klaster. Mund të përmirësohet situata me përpunim të parakohshëm. Vendosa të provoj të ndaj të dhënat e papërpunuara në kolona kromozomesh, kështu që shpresoja t'i ofroja Spark-u të dhënat "të parakohshme të ndara".

Kërkova në StackOverflow se si të ndaj sipas vlerave të kolonave dhe gjeta një përgjigje kaq të bukur. Me AWK, mund të ndani një skedar teksti sipas vlerave të kolonave, duke bërë një regjistrim në skript, në vend që të dërgoni rezultatet në stdout.

Për prova, shkrova një skript Bash. Shkova një nga paketat e TSV, pastaj e dekompresova me gzip dhe e dërgova në awk.

gzip -dc path/to/chunk/file.gz |
awk -F 't' 
'{print $1",..."$30">"chunked/"$chr"_chr"$15".csv"}'

Kjo funksionoi!

Plotësimi i bërthamave

ÇfarĂ« mĂ«sova: gnu parallel Ă«shtĂ« diçka magjike, çdo kush duhet ta pĂ«rdorĂ« atĂ«.

Ndarja po shkonte mjaft ngadalë, dhe kur e nisa htop, për të kontrolluar përdorimin e instancës EC2 (dhe të shtrenjtë) të fuqishme, u zbulua se po përdorja vetëm një bërthamë dhe rreth 200 MB memorie. Për të zgjidhur problemin dhe për të mos humbur shumë para, duhej të mendoja se si të parallelizoja punën. Me fat, në një libër të mrekullueshëm Data Science at the Command Line të Jeron Jansens, gjeta një kapitull që merret me parallelizimin. Nga ai mësova për gnu parallel, një metodë shumë fleksibël për implementimin e multi-threading në Unix.

Parse 25TB me AWK dhe R
Kur e nisa ndarjen me procesin e ri, gjithçka ishte shkëlqyer, por kishte një ngushticë - shkarkimi i objekteve S3 në disk nuk ishte shumë i shpejtë dhe nuk ishte plotësisht i parallelizuar. Për ta korrigjuar këtë, bëra këto:

  1. Pashë se mund të realizoja fazën e shkarkimit S3 direkt nëpërmjet tubit, duke përjashtuar plotësisht ruajtjen përkohësisht në disk. Kjo do të thotë se mund të shmangia shkrimin e të dhënave të papërpunuara në disk dhe të përdorja një ruajtje edhe më të vogël, dhe kështu më të lira në AWS.
  2. Me komandën aws configure set default.s3.max_concurrent_requests 50 rrit shumë numrin e threads që përdor AWS CLI (në mënyrë standarde janë 10).
  3. Kalova në një instancë EC2 të optimizuar për shpejtësinë e rrjetit, me shkronjën n në emër. Zbulova se humbja e fuqisë llogaritëse duke përdorur instancat n u kompensua me shpejtësinë e ngarkesës. Për shumicën e detyrave përdora c5n.4xl.
  4. Ndryshova gzip në pigz, kjo është një mjet gzip që di të bëjë gjëra të mrekullueshme për të parallelizuar një detyrë fillimisht jo të parallelizuar të dekompresimit të skedarëve (kjo ndihmoi më pak).

# Let S3 use as many threads as it wants
aws configure set default.s3.max_concurrent_requests 50

for chunk_file in $(aws s3 ls $DATA_LOC | awk '{print $4}' | grep 'chr'$DESIRED_CHR'.csv') ; do

        aws s3 cp s3://$batch_loc$chunk_file - |
        pigz -dc |
        parallel --block 100M --pipe  
        "awk -F 't' '{print $1",..."$30">"chunked/{#}_chr"$15".csv"}'"

       # Combine all the parallel process chunks to single files
        ls chunked/ |
        cut -d '_' -f 2 |
        sort -u |
        parallel 'cat chunked/*_{} | sort -k5 -n -S 80% -t, | aws s3 cp - '$s3_dest'/batch_'$batch_num'_{}'
        
         # Clean up intermediate data
       rm chunked/*
done

Këto hapa u kombinuan me njëri-tjetrin që të gjitha të funksiononin shumë shpejt. Falë rritjes së shpejtësisë së shkarkimit dhe heqjes së shkrimit në disk, tani mund të procesoj një paketë 5 terabajtësh për vetëm disa orë.

Ky tweet do tĂ« duhej tĂ« pĂ«rmendte ‘TSV’. fatkeqĂ«sisht.

Përdorimi i të dhënave të riparasura

ÇfarĂ« mĂ«sova: Spark-i do tĂ« dhĂ«na tĂ« pambushura dhe nuk e do tĂ« kombinojĂ« partitĂ«.

Tani tĂ« dhĂ«nat ishin nĂ« S3 nĂ« njĂ« format tĂ« papackuar (lexo, tĂ« ndara) dhe gjysmĂ« tĂ« renditur, dhe mund tĂ« kthehesha pĂ«rsĂ«ri te Spark. MĂ« priste njĂ« surprizĂ«: sĂ«rish nuk arrija tĂ« arrija atĂ« qĂ« dĂ«shiroja! Ishte shumĂ« e vĂ«shtirĂ« t’i thoja Spark-ut se si ishin tĂ« partizuara tĂ« dhĂ«nat. Dhe edhe kur e bĂ«ra kĂ«tĂ«, rezultoi se kishim shumĂ« patica (95 mijĂ«), dhe kur me pomocĂ«n e coalesce e zvogĂ«lova atĂ« nĂ« kufij tĂ« arsyeshĂ«m, prishja rregullimin tim tĂ« particioneve. Jam i sigurt se mund tĂ« rregullohet, por pas disa ditĂ«sh kĂ«rkime nuk arrita tĂ« gjej njĂ« zgjidhje. NĂ« fund, e plotĂ«sova tĂ« gjitha detyrat nĂ« Spark, edhe pse kjo mori pak kohĂ«, dhe skedarĂ«t e mi tĂ« ndarĂ« Parquet nuk ishin shumĂ« tĂ« vegjĂ«l (~200 KB). MegjithatĂ«, tĂ« dhĂ«nat ishin aty ku duhej.

Parse 25TB me AWK dhe R
Shumë të vogla dhe të ndryshme, mrekulli!

Testimi i kërkesave lokale në Spark

ÇfarĂ« mĂ«sova: nĂ« Spark ka shumĂ« tĂ« meta kur zgjidhni detyra tĂ« thjeshta.

Duke ngarkuar të dhënat në një format të menduar, arrita të testoj shpejtësinë. Përshtata skriptin në R për të nisur një server lokal Spark, dhe më pas ngarkova kornizën e të dhënave Spark nga magazina e caktuar Parquet-bins. Përpika të ngarkoja të gjithë të dhënat, por nuk arrita ta bëj Sparklyr të njohë partizionimin.

sc <- Spark_connect(master = "local")

desired_snp <- 'rs34771739'

# Filloni një kohëmatës
start_time <- Sys.time()

# Ngarko binin e dëshiruar në Spark
intensity_data % 
  Spark_read_Parquet(
    name = 'intensity_data', 
    path = get_snp_location(desired_snp),
    memory = FALSE )

# Ndarja e binit në snp dhe pastaj mbledhja në lokal
test_subset % 
  filter(SNP_Name == desired_snp) %>% 
  collect()

print(Sys.time() - start_time)

Ekzekutimi zgjati 29.415 sekonda. Shumë më mirë, por jo aq mirë për testimin masiv të diçkaje. Për më tepër, nuk mund të shkurtova kohën me ndihmën e caching, sepse kur përpiqesha të ruaja në memory kornizën e të dhënave, Spark gjithmonë dështonte, edhe kur ndaja më shumë se 50 GB memorie për një dataset që peshoi më pak se 15.

Kthimi te AWK

ÇfarĂ« mĂ«sova: strukturat associative nĂ« AWK janĂ« shumĂ« efikase.

E dija se mund të arrija shpejtësi më të lartë. Më kujtohej një herë kur isha duke lexuar një udhëzues të shkëlqyer për AWK nga Bruce Barnett Kam lexova për një karakteristikë të shkëlqyer që quhet "marrëdhënie asociative". Në thelb, janë çifte çelësi-vlerë, që për një arsye e ndryshuan emrin në AWK dhe për këtë arsye nuk i kisha kujtuar shumë. Roman Cheplyaka më kujtoi se termi "marrëdhënie asociative" është shumë më i vjetër se termi "çelës-vlerë". Edhe nëse ju kërkoni çelës-vlerë në Google Ngram, nuk do ta gjeni këtë term, por do të gjeni marrëdhënie asociative! Për më tepër, "çelësi-vlerë" shpesh lidhet me bazat e dhënave, prandaj është shumë më logjike të krahasohet me hashmap. E kuptova se mund t'i përdor këto marrëdhënie asociative për të lidhur SNP-të e mia me tabelën e grupeve (tabela binare) dhe të dhënat e papërpunuara pa përdorur Spark.

Për këtë, në skenarin AWK përdora bllokun FILLON. Ky është një fragment kodi që ekzekutohet para se të kalojë rreshti i parë i të dhënave në trupin kryesor të skenarit.

join_data.awk
BEGIN {
  FS=",";
  batch_num=substr(chunk,7,1);
  chunk_id=substr(chunk,15,2);
  while(getline  "chunked/chr_"chr"_bin_"bin[$1]"_"batch_num"_"chunk_id".csv"
}

Ekipa while(getline...) ngarkohet të gjitha rreshtat nga CSV e grupit (binar), caktohet kolonën e parë (emri i SNP) si çelësi për marrëdhënien asociative bin dhe vlera e dytë (grupi) si vlera. Pastaj në bllokun { }, i cili ekzekutohet për të gjitha rreshtat e skedarit kryesor, çdo rresht dërgohet në skedarin e daljes, që merr një emër unik në varësi të grupit të tij (binar): ..._bin_"bin[$1]"_....

Variablat batch_num dhe chunk_id përputheshin me të dhënat e ofruara nga tubi, duke shmangur kështu një gjendje konkurruese, dhe çdo rrjedhë ekzekutimi e nisur paralelet, shkruante në skedarin e saj unik.

Dukeqë të gjitha të dhënat e papërpunuara i kisha shpërndarë në dosje sipas kromosomeve, që kishin mbetur pas eksperimenteve të mia të mëparshme me AWK, tani mund të shkruaja një skenar tjetër Bash për të përpunuar një kromozom në një herë dhe dorëzuar në S3 të dhëna më të ndara.

DESIRED_CHR='13'

# Shkarko të dhënat e kromozomeve nga s3 dhe ndaj në bins
aws s3 ls $DATA_LOC |
awk '{print $4}' |
grep 'chr'$DESIRED_CHR'.csv' |
parallel "echo 'po lexoj {}'; aws s3 cp "$DATA_LOC"{} - | awk -v chr=""$DESIRED_CHR"" -v chunk="{}" -f split_on_chr_bin.awk"

# Bashko të gjitha proceset paralele në skedarë të vetëdijshëm dhe ngarko në rds duke përdorur R
ls chunked/ |
cut -d '_' -f 4 |
sort -u |
parallel "echo 'po zip' { }; cat chunked/*_bin_{}_*.csv | ./upload_as_rds.R '$S3_DEST'/chr_'$DESIRED_CHR'_bin_{}.rds"
rm chunked/*

Në skenar ka dy pjesë paralelet.

Në seksionin e parë, të dhënat lexohen nga të gjitha skedarët që përmbajnë informacionin për kromozomën e nevojshme, pastaj këto të dhëna shpërndahen në thread-e, të cilat i ndajnë skedarët në grupet përkatëse (bin). Për të shmangur një gjendje garash, kur disa thread-e shkruajnë në një skedar të vetëm, AWK kalon emrat e skedarëve për të shkruar të dhënat në vende të ndryshme, për shembull, chr_10_bin_52_batch_2_aa.csv. Si rezultat, krijohet një numër i madh skedash të vegjël në disk (për këtë kam përdorur vëllime EBS me kapacitet një terabajt).

Përçuesi nga seksioni i dytë paralelet kalon nëpër grupet (bin) dhe bashkon skedarët e tyre të veçantë në CSV të përbashkët cat, dhe pastaj i dërgon ata për eksport.

Translacioni në R?

ÇfarĂ« mĂ«sova: mund tĂ« referohesh nĂ« stdin dhe stdout nga skripti R, çka do tĂ« thotĂ« qĂ« mund ta pĂ«rdorĂ«sh atĂ« nĂ« pĂ«rçues.

Në skriptin Bash, mund të keni vënë re këtë rresht: ...cat chunked/*_bin_{}_*.csv | ./upload_as_rds.R.... Ai transfton të gjithë skedarët e lidhur të grupit (bin) në skriptin R më poshtë. {} është një metodë speciale paralelet, e cila çdo të dhënë të dërguar në atë rrjedhë e vendos direkt në komandën e vet. Opsioni {#} siguron një ID unik të rrjedhës së ekzekutimit, dhe {%} përfaqëson numrin e slotit të detyrës (ripërseritet, por kurrë në të njëjtën kohë). Lista e të gjitha opsioneve mund të gjendet në të dokumentacionit.

#!/usr/bin/env Rscript
library(readr)
library(aws.s3)

# Read first command line argument
data_destination <- commandArgs(trailingOnly = TRUE)[1]

data_cols <- list(SNP_Name = 'c', ...)

s3saveRDS(
  read_csv(
        file("stdin"), 
        col_names = names(data_cols),
        col_types = data_cols 
    ),
  object = data_destination
)

Kur variabla file("stdin") i kalon në readr::read_csv, të dhënat e transfuara në skriptin R ngarkohen në një frame, i cili më pas në formën e .rds-skedarit me anë të aws.s3 shkruhet direkt në S3.

RDS është ngjashëm me një version më të vogël të Parquet, pa rafinimet e magazinimit kolonor.

Pas përfundimit të skriptit Bash, kam marrë një grumbull .rds-skedash në S3, që më lejoj të përdor kompresim efektiv dhe lloje të integrate.

Pavarësisht nga përdorimi ngadalësues të R, gjithçka punonte shumë shpejt. Nuk është e habitshme që fragmentet në R, që janë përgjegjëse për të lexuar dhe shkruar të dhëna, janë optimizuar mirë. Pas testimit në një kromozom me madhësi mesatare, detyra u përfundua në instancën C5n.4xl për rreth dy orë.

Kufizimet S3

ÇfarĂ« mĂ«sova: falĂ« implementimit tĂ« mençur tĂ« rrugĂ«ve, S3 mund tĂ« pĂ«rpunojĂ« njĂ« numĂ«r tĂ« madh skedash.

Ishte shqetësim nëse S3 do të mund të procesonte shumicën e skedarëve që i dërgoheshin. Mund të kisha bërë emrat e skedarëve kuptimplotë, por si do të kërkonte S3 sipas tyre?

Parse 25TB me AWK dhe R
Folderët në S3 janë thjesht për bukuri, në të vërtetë sistemi nuk i intereson simboli /. Nga faqja FAQ e S3.

Duket se S3 paraqet rrugën deri te një skedar të caktuar si një çelës të thjeshtë në një tabelë hash ose në një bazë të dhënash të bazuar në dokumente. Cisternat (buckets) mund të merren si tabela, ndërsa skedarët si regjistrime në këtë tabelë.

Pasi që shpejtësia dhe efikasiteti janë të rëndësishme për realizimin e fitimeve në Amazon, nuk është çudi që ky sistem "çelës-në-si-rrugë-deri-te-skedarin" është jashtëzakonisht i optimizuar. Kam tentuar të gjej një ekuilibër: që të mos jetë e nevojshme të bëhen shumë kërkesa get, por që kërkesat të ekzekutohen shpejt. Doli që më së miri të bëhen rreth 20,000 skedarë binarë. Mendoj se nëse vazhdoj të optimizoj, mund të arrij një rritje të shpejtësisë (për shembull, të bëj një cisternë të veçantë vetëm për të dhënat, duke zvogëluar kështu madhësinë e tabelës së kërkimit). Megjithatë, nuk kishte më kohë dhe mundësi për eksperimente të mëtejshme.

ÇfarĂ« ndodh me pĂ«rputhshmĂ«rinĂ« ndĂ«rkufitare?

ÇfarĂ« mĂ«sova: arsyeja kryesore e humbjes sĂ« kohĂ«s Ă«shtĂ« optimizimi i parakohshĂ«m i metodĂ«s suaj tĂ« ruajtjes.

Në këtë moment është shumë e rëndësishme të pyesni veten: "Pse të përdorim një format skedari pronësor?" Arsyeja qëndron në shpejtësinë e ngarkimit (skedarët gzip CSV të paketuar ngarkoheshin 7 herë më gjatë) dhe përputhshmërinë me proceset tona të punës. Mund të rishikoj vendimin tim nëse R do të ishte në gjendje të ngarkonte lehtësisht skedarët Parquet (apo Arrow) pa ngarkesa nga Spark. Të gjithë në laboratorin tonë përdorin R, dhe nëse duhet të konvertoj të dhënat në një format tjetër, kam ende të dhënat origjinale tekstuale, kështu që mund ta filloj përsëri procesin.

Ndara e punës

ÇfarĂ« mĂ«sova: mos u pĂ«rpoqni tĂ« optimizoni detyrat manualisht, lejoje kĂ«tĂ« ta bĂ«jĂ« kompjuteri.

Kam debatuar procesin e punës në një kromozom, tani është koha të përpunoj të dhënat e tjera.
Dëshiroja të ngrija disa instancë EC2 për konvertim, por në të njëjtën kohë kisha frikë se do të merrja ngarkesë jashtëzakonisht të paekuilibruar në detyrat e përpunimit (ashtu si Spark kishte vuajtur nga partitë e paekuilibruar). Për më tepër, nuk më pëlqente të ngrija një instancë për çdo kromozom, sepse ka një kufizim të paracaktuar prej 10 instancash për llogaritë AWS.

Atëherë vendosa të shkruaj një skritp në R për optimizimin e detyrave të përpunimit.

Së pari kërkova nga S3 të llogariste sesa hapësirë në ruajtje zë çdo kromozom.

library(aws.s3)
library(tidyverse)

chr_sizes % 
  mutate(Size = as.numeric(Size)) %>% 
  filter(Size != 0) %>% 
  mutate(
    # Nxjerr kromozomin nga emri i skedarit 
    chr = str_extract(Key, 'chr.{1,4}.csv') %>%
             str_remove_all('chr|.csv')
  ) %>% 
  group_by(chr) %>% 
  summarise(total_size = sum(Size)/1e+9) # Ndaj për të marrë vlerën në GB



# Një tibble: 27 x 2
   chr   total_size
         
 1 0           163.
 2 1           967.
 3 10          541.
 4 11          611.
 5 12          542.
 6 13          364.
 7 14          375.
 8 15          372.
 9 16          434.
10 17          443.
# 
 me 17 rreshta të tjerë

Më pas shkrova një funksion që merr dimensionin total, përzjen renditjen e kromozomeve, i ndan në grupe. num_jobs dhe raporton se sa ndryshojnë dimensionet e të gjitha detyrave për procesim.

num_jobs <- 7
# Sa i madh do të ishte secili punë nëse do të ndaheshin perfekt?
job_size <- sum(chr_sizes$total_size)/7

shuffle_job %
    sample_frac() %>% 
    mutate(
      cum_size = cumsum(total_size),
      job_num = ceiling(cum_size/job_size)
    ) %>% 
    group_by(job_num) %>% 
    summarise(
      job_chrs = paste(chr, collapse = ','),
      total_job_size = sum(total_size)
    ) %>% 
    mutate(sd = sd(total_job_size)) %>% 
    nest(-sd)
}

shuffle_job(1)



# Një tibble: 1 x 2
     sd data            
             
1  153.

Më pas e kam drejtuar me purrr një mijë përzierje dhe kam zgjedhur më të mirën.

1:1000 %>% 
  map_df(shuffle_job) %>% 
  filter(sd == min(sd)) %>% 
  pull(data) %>% 
  pluck(1)

Kështu mora një set detyrash, shumë të ngjashme në madhësi. Më pas mbeti vetëm ta vështjelloj skriptin tim të mëparshëm në një cikël të madh. për. Shkrimi i kësaj optimizimi zuri rreth 10 minuta. Dhe kjo është shumë më pak sesa do të kisha shpenzuar për krijimin manual të detyrave në rast të pamjaftueshmërisë së tyre. Prandaj mendoj se me këtë optimizim paraprak nuk kam gabuar.

për DESIRED_CHR në "16" "9" "7" "21" "MT"
do
# Kodi për procesimin e një kromozomi të vetëm
fi

Në fund shtoj komandën e fikjes:

sudo shutdown -h now


 dhe gjithçka funksionoi! Me ndihmĂ«n e AWS CLI ngrita instanca dhe pĂ«rmes opsionit user_data i kalova skriptet e Bash-eve tĂ« detyrave pĂ«r procesim. Ato u ekzekutuan dhe u fikĂ«n automatikisht, kĂ«shtu qĂ« nuk paguaj pĂ«r forcĂ«n kompjuterike tĂ« tepĂ«rt.

aws ec2 run-instances ...
--tag-specifications "ResourceType=instance,Tags=[{Key=Name,Value=<>}]" 
--user-data file://<>

Po paketojmë!

ÇfarĂ« mĂ«sova: API duhet tĂ« jetĂ« i thjeshtĂ« pĂ«r shkak tĂ« thjeshtĂ«sisĂ« dhe fleksibilitetit tĂ« pĂ«rdorimit.

Më në fund mora të dhënat në vendin dhe formën e dëshiruar. Duhej të thjeshtoja procesin e përdorimit të të dhënave, që kolegët e mi të kishin më të lehtë. Doja të bëja një API të thjeshtë për krijimin e kërkesave. Nëse në të ardhmen vendos të kaloj nga .rds në skedarët Parquet, kjo duhet të jetë një problem për mua, jo për kolegët. Prandaj, vendosa të krijoj një paketë të brendshme R.

Mblodha dhe dokumentova një paketë shumë të thjeshtë, e cila përmban vetëm disa funksione për qasjen në të dhënat, të mbledhura rreth funksionit get_snp. Gjithashtu, për kolegët krijova një faqe pkgdown, në mënyrë që ata të mund ta shohin lehtësisht shembujt dhe dokumentacionin.

Parse 25TB me AWK dhe R

Keqja inteligjente

ÇfarĂ« mĂ«sova: nĂ«se tĂ« dhĂ«nat tuaja janĂ« mirĂ« pĂ«rgatitur, do tĂ« jetĂ« e lehtĂ« tĂ« keqkuptohet!

Duke qenë se një nga proceset kryesore aplikonte të njëjtin model analize në paketën SNP, vendosa ta përdor grupimin (binning) në përfitimin tim. Kur kalojmë të dhënat nga SNP, informacioni nga grupa (bin) bashkëngjitet me objektin e kthyer. Kështu, kërkesat e vjetra mund (teoretikisht) të shpejtojnë përpunimin e kërkesave të reja.

# Part of get_snp()
...
  # Test if our current snp data has the desired snp.
  already_have_snp <- desired_snp %in% prev_snp_results$snps_in_bin

  if(!already_have_snp){
    # Grab info on the bin of the desired snp
    snp_results <- get_snp_bin(desired_snp)

    # Download the snp's bin data
    snp_results$bin_data <- aws.s3::s3readRDS(object = snp_results$data_loc)
  } else {
    # The previous snp data contained the right bin so just use it
    snp_results <- prev_snp_results
  }
...

Gjatë ndërtimit të paketës, bëra shumë benchmark-e për të krahasuar shpejtësinë e përdorimit të metodave të ndryshme. E rekomandoj atë, sepse ndonjëherë rezultatet janë të papritura. Për shembull, dplyr::filter u tregua shumë më i shpejtë se sa kapja e rreshtave me filtrimin në bazë të indeksimit, dhe marrja e një kolone nga korniza e dhënash e filtruar punonte shumë më shpejt se sa përdorimi i sintaksës së indeksimit.

Kujdes, që objekti prev_snp_results përmban çelësin snps_in_bin. Ky është një array i të gjitha SNP-ve unike në grup (bin), duke lejuar kontrollin e shpejtë për të parë nëse të dhënat nga kërkesa e mëparshme tashmë ekzistojnë. Gjithashtu, e thjeshton kalimin përmes të gjitha SNP-ve në grup (bin) me këtë kod:

# Get bin-mates
snps_in_bin <- my_snp_results$snps_in_bin

for(current_snp in snps_in_bin){
  my_snp_results <- get_snp(current_snp, my_snp_results)
  # Do something with results 
}

Rezultatet

Tani mund të (dhe filluam seriozisht) të përjashtojmë modele dhe skenarë, të cilat më pare ishin të paarritshme për ne. Më e mira është se kolegët e mi në laborator nuk duhet të mendojnë për ndonjë ndërlikim. Ata thjesht kanë një funksion që funksionon.

Dhe ndonëse paketa i çliron ata nga detajet, përpiqesha ta bëja formatin e të dhënave të mjaftueshëm të thjeshtë, në mënyrë që ata të mund të kuptonin nëse nesër unë papritur zhdukem...

Shpejtësia është rritur ndjeshëm. Zakonisht ne skanojmë fragmente funksionale të rëndësishme të genomit. Më parë nuk mund të bënim këtë (ishte shumë e shtrenjtë), por tani, falë strukturës grupore (bin) dhe keqkuptimit, kërkesa për një SNP mesatarisht zgjat më pak se 0,1 sekondë, dhe përdorimi i të dhënave është aq i ulët, sa shpenzimet për S3 janë shumë të vogla.

Përfundim

Ky artikull nuk është një udhëzues fare. Zgjidhja rezultoi të ishte individuale dhe pothuajse me siguri jo optimale. Në të vërtetë, ky është një tregim për një udhëtim. Dua që të tjerët të kuptojnë se zgjidhjet e tilla nuk ndodhin plotësisht të formuara në mendje, ato janë rezultat i provave dhe gabimeve. Për më tepër, nëse jeni në kërkim të një specialisti në analizën e të dhënave, merrni parasysh se për përdorimin efektiv të këtyre instrumenteve nevojitet përvojë, dhe përvoja kërkon para. Jam i lumtur që pata fonde për të paguar, por shumë të tjerë, që mund të bëjnë punën e njëjtë më mirë se unë, kurrë nuk do të kene këtë mundësi për shkak të mungesës së parave edhe për të provuar.

Mjetet për të dhëna të mëdha janë universale. Nëse keni kohë, pothuajse me siguri do të jeni në gjendje të shkruani një zgjidhje më të shpejtë duke aplikuar "pastrimin inteligjent" të të dhënave, ruajtjen dhe metodologjitë e nxjerrjes. Në fund të fundit, gjithçka vjen në analizën e shpenzimeve dhe përfitimeve.

ÇfarĂ« kam mĂ«suar:

  • nuk ka njĂ« mĂ«nyrĂ« tĂ« lirĂ« pĂ«r tĂ« procesuar 25 TB tĂ« dhĂ«na njĂ«herĂ«sh;
  • tĂ« jeni tĂ« kujdesshĂ«m me madhĂ«sinĂ« e skedarĂ«ve tuaj Parquet dhe organizimin e tyre;
  • particionet nĂ« Spark duhet tĂ« jenĂ« tĂ« balancuara;
  • kurrĂ« mos pĂ«rpiquni tĂ« krijoni 2.5 milion partitore;
  • grupimi mbetet i vĂ«shtirĂ«, ashtu si konfigurimi i Spark;
  • ndonjĂ«herĂ« tĂ« dhĂ«nat e veçanta kĂ«rkojnĂ« zgjidhje tĂ« veçanta;
  • bashkimi nĂ« Spark funksionon shpejt, por partitoret ende kanĂ« kosto tĂ« larta;
  • mos flini kur ju mĂ«sojnĂ« bazat, sigurisht qĂ« dikush tashmĂ« e ka zgjidhur problemin tuaj qĂ« nĂ« vitet 1980;
  • gnu parallel - Ă«shtĂ« njĂ« gjĂ« magjike, tĂ« gjithĂ« duhet ta pĂ«rdorin;
  • Spark preferon tĂ« dhĂ«na tĂ« pakompresuara dhe nuk dĂ«shiron tĂ« kombinojĂ« partitoret;
  • nĂ« Spark ka shumĂ« overhead pĂ«r zgjidhjen e problemeve tĂ« thjeshta;
  • çfarĂ« grupimeve associative nĂ« AWK janĂ« shumĂ« efikase;
  • mund tĂ« adresohet stdin dhe stdout nga skripti R, kĂ«shtu qĂ« mund ta pĂ«rdorĂ«sh atĂ« nĂ« pipeline;
  • falĂ« implementimit inteligjent tĂ« ndalesave S3 mund tĂ« pĂ«rpunojĂ« shumĂ« skedarĂ«;
  • arsyeja kryesore e humbjes sĂ« kohĂ«s Ă«shtĂ« optimizimi i parakohshĂ«m i metodĂ«s suaj tĂ« ruajtjes;
  • mos u pĂ«rpiqni tĂ« optimizoni detyrat manualisht, lejoni kompjuterin ta bĂ«jĂ« atĂ«;
  • API duhet tĂ« jetĂ« i thjeshtĂ« pĂ«r thjeshtĂ«sinĂ« dhe fleksibilitetin e pĂ«rdorimit;
  • nĂ«se tĂ« dhĂ«nat tuaja janĂ« tĂ« pĂ«rgatitura mirĂ«, cache do tĂ« jetĂ« i lehtĂ«!

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster