Parsing 25TB using AWK and R

Parsing 25TB using AWK and R
Kuidas seda artiklit lugeda: palun vabandust, et tekst lĂ€ks nii pikaks ja segaseks. Ajan teie aega sÀÀstmiseks alustasin iga peatĂŒki sissejuhatusega „Mida ma Ă”ppisin”, kus ma lĂŒhidalt selgitĂ€n peatĂŒki sisu ĂŒhe vĂ”i kahe lausega.

„Lihtsalt nĂ€ita lahendust!” Kui soovite lihtsalt nĂ€ha, kuhu ma jĂ”udsin, minge peatĂŒkki „Ajan end leidlikumaks”, kuid ma arvan, et huvitavam ja kasulikum on lugeda ebaĂ”nnestumistest.

Hiljuti anti mulle ĂŒlesanne seadistada protsess suure koguse DNA algse jĂ€rjestuse (tehniliselt SNP-chip) töötlemiseks. Selleks pidi kiiresti saama teavet antud geneetilise asukoha (mida kutsutakse SNP-ks) kohta, et hiljem saaks tegeleda modelleerimise ja muude ĂŒlesannetega. Kasutades R ja AWK, sain andmeid loogiliselt puhastada ja korraldada, mis kiirendas oluliselt pĂ€ringute töötlemist. See ei olnud lihtne ja nĂ”udis arvukaid iteratsioone. See artikkel aitab teil vĂ€ltida mĂ”ningaid minu vigu ja demonstreerib, millega ma lĂ”puks hakkama sain.

Alustuseks mÔned tutvustavad selgitused.

Andmed

Meie ĂŒlikooli geneetilise informatsiooni töötlemise keskus andis meile andmeid TSV vormingus, kogus 25 TB. Need olid jagatud 5 paketti, Gzip-i tihendatud, kus igaĂŒhes oli umbes 240 nelja gigabaidi faili. Iga rida sisaldas andmeid ĂŒhe SNP kohta ĂŒhe isiku kohta. Kokku edastati andmed umbes 2,5 miljoni SNP ja umbes 60 tuhande inimese kohta. Peale SNP-teabe sisaldasid failid palju veerge numbritega, mis kajastasid erinevaid omadusi, nagu lugemise intensiivsus, erinevate alleelide sagedus jne. Kokku oli umbes 30 veergu unikaalsete vÀÀrtustega.

EesmÀrk

Nagu igas andmehĂ€lgema projektis, oli kĂ”ige olulisem mÀÀrata, kuidas andmeid kasutatakse. Antud juhul peame enamasti valima mudeleid ja töövooge SNP-de pĂ”hjal SNP-de jĂ€rgi. See tĂ€hendab, et samal ajal on meil vaja andmeid ainult ĂŒhe SNP kohta. Ma pidin Ă”ppima, kuidas kĂ”ige lihtsamalt, kiiremini ja odavamalt vĂ€lja vĂ”tta kĂ”ik kirjed, mis on seotud ĂŒhe 2,5 miljoni SNP-st.

Kuidas seda mitte teha

Tsiteerin sobivat kliĆĄeed:

Ma ei ole tuhat korda ebaÔnnestunud, vaid avastasin tuhat viisi, kuidas mitte andmeid hÔlpsaks pÀringuks parseerida.

Esimene katse

Mida ma Ôppisin: ei ole odavat viisi 25 TB korraga eemaldada.

Kuulates Vanderbilti Ülikooli kursust "Suured andmed: arenenud töötlemismeetodid", olin kindel, et see on lihtne. KĂŒllap kulub tund vĂ”i kaks Hive-serveri seadistamiseks, et kĂ”ik andmed lĂ€bi kĂ€ia ja tulemusi esitada. Kuna meie andmed asuvad AWS S3-s, kasutasin teenust Athena, mis vĂ”imaldab rakendada Hive SQL-pĂ€ringuid S3 andmetele. Hive-kluster ei vaja seadistamist/ĂŒlesehitamist ja maksad vaid nende andmete eest, mida otsid.

PÀrast seda, kui nÀitasin Athenale oma andmeid ja nende vormingut, jooksin mÔned testid sarnaste pÀringutega:

select * from intensityData limit 10;

Ja sain kiiresti korralikult struktureeritud tulemusi. Valmis.

Kuni pole proovinud andmeid töös


Paluti mul vÀlja vÔtta kogu info SNP-de kohta, et testida oma mudelit. KÀivitasin pÀringu:


select * from intensityData 
where snp = 'rs123456';


ja hakkasin ootama. Kaheksa minuti ja ĂŒle 4 TB kĂŒsitud andmete jĂ€rel sain tulemuse. Athena tasub leitud andmete mahu eest, $5 terabaidi kohta. Seega maksis see ainus pĂ€ring $20 ja ooteaega kaheksa minutit. Mudeli lĂ€biviimiseks kĂ”igi andmete pĂ”hjal oleks pidanud ootama 38 aastat ja maksma $50 miljonit. Ilmselgelt ei sobinud see meile.

Pidiime kasutama Parquet


Mida ma Ôppisin: olge ettevaatlikud oma Parquet-failide suuruse ja korraldusega.

Esialgu proovisin olukorda parandada, konverteerides kÔik TSV-d Parquet-failideks. Need on mugavad suure andmehulga töötlemiseks, sest teave on salvestatud veergude kaupa: iga veerg asub enda mÀlusegmendis/diskis, erinevalt tekstifailidest, kus read sisaldavad iga veeru elemente. Ja kui midagi otsida, piisab vajalikust veerust lugemisest. Lisaks sisaldab iga faili veerus vÀÀrtuste vahemik, nii et kui otsitav vÀÀrtus ei ole veeru vahemikus, ei raiska Spark aega kogu faili skaneerimisele.

KĂ€ivitasin lihtsa ĂŒlesande AWS Glue kuidas meie TSV-d Parquet-ks konverteerida ja uued failid Athenasse ĂŒles laadida. Selleks kulus umbes 5 tundi. Kuid kui ma kĂ€ivitasin pĂ€ringu, siis selle tĂ€itmiseks kulus umbes sama palju aega ja veidi vĂ€hem raha. Asi on selles, et Spark, proovides ĂŒlesannet optimeerida, lihtsalt pakkis ĂŒhe TSV-osa lahti ja pani selle oma Parquet-osa. Ja kuna iga osa oli piisavalt suur ja sisaldas tervet rida paljude inimeste registreeringutest, siis sisaldas iga fail kĂ”iki SNP-sid, seetĂ”ttu pidi Spark avama kĂ”ik failid, et otsida vajalikku teavet.

Huvitav on see, et Parquet'i vaikimisi (ja soovitatav) kompressioonitĂŒĂŒp — snappy — ei ole jagatav (splitable). SeetĂ”ttu takerdus iga tĂ€itja (executor) ĂŒlesande lahti pakkimise ja kogu 3,5 GB suuruse andmesessiooni ĂŒleslaadimise peale.

Parsing 25TB using AWK and R

KÀime probleemiga lÀbi

Mida ma Ôppisin: sortimine on keeruline, eriti kui andmed on hajutatud.

Mul nĂ€is, et nĂŒĂŒd sain probleemist aru. Peasin lihtsalt andmed sorteerima SNP veeru alusel, mitte inimeste alusel. Siis sĂ€ilitatakse eraldi andmeosas mitu SNP-d ja siis tuleb pargu „nutikas“ funktsioon mĂ€ngu, mis ĂŒtleb „ava ainult siis, kui vÀÀrtus on vahemikus“. Kahjuks osutus miljardeid ridu, mis olid hajutatud klastrisse, sorteerimine keerukaks ĂŒlesandeks.

AWS ei taha kindlasti raha tagasi maksta pĂ”hjendusega „Ma olen hajevil ĂŒliĂ”pilane“. PĂ€rast seda, kui kĂ€ivitasin sortimise Amazon Glue'is, töötas see 2 pĂ€eva ja lĂ”petas ebaĂ”nnestumisega.

Kuidas on partitionimist?

Mida ma Ôppisin: Sparkis peavad osad olema tasakaalustatud.

Siis tuli mulle idee andmeid kromosoomide jÀrgi partitionsida. Neid on 23 (ja veel mÔned, kui arvestada mitoohondria DNA-d ja dekodeerimata alasid).
See vĂ”imaldab andmed jagada vĂ€iksemateks osadeks. Kui lisada Spark-funktsiooni ekspordi skripti Glue'ile lihtsalt ĂŒks rida partition_by = "chr", siis andmed peaksid olema jagatud kĂ€rudesse (buckets).

Parsing 25TB using AWK and R
Genoom koosneb arvukatest fragmentidest, mida nimetatakse kromosoomideks.

Kahjuks see ei toiminud. Chromosoomidel on erinevad suurused, mis tĂ€hendab ka erinevat teabe hulka. See tĂ€hendab, et Spark'i kaudu töötajatele saadetud ĂŒlesanded ei olnud tasakaalus ja tĂ€ideti aeglaselt, kuna mĂ”ned sĂ”lmed lĂ”petasid varem ja seisisid. Kuid ĂŒlesanded viidi ellu. Siiski, kui kĂŒsiti ĂŒht SNP, taasiseseisvus tasakaalustamatus probleemideks. SNP töötlemise maksumus suuremates kromosoomides (st sealt, kust me andmeid soovime saada) vĂ€henes vaid umbes 10 korda. Palju, aga mitte piisavalt.

Ent kui jagada veel vÀiksemateks partiideks?

Mida ma Ă”ppisin: ĂŒldiselt Ă€rge proovige kunagi teha 2,5 miljonit partiid.

Otsustasin kĂ”ike proovida ja partitseerisin iga SNP. See tagas, et partii suurused oleksid ĂŒhtlased. SEE OLI HALB IDEE. Kasutasin Glue'i ja lisasin sĂŒĂŒtu rea partition_by = 'snp'. Ülesanne kĂ€ivitati ja hakkas jooksma. PĂ€eva pĂ€rast kontrollisin ja nĂ€gin, et S3-s pole ikka veel midagi salvestatud, nii et lĂ”petasin ĂŒlesande. Tundub, et Glue salvestas vahefailid S3-sse varjatud kohta, ja palju faile, vĂ”ib-olla paar miljonit. Selle tulemusena maksis minu viga rohkem kui tuhat dollarit ja ei rÔÔmustanud mu mentorit.

Partitsioneerimine + sortimine

Mida ma Ôppisin: sortimine on endiselt keeruline, nagu ka Spark'i seadistamine.

Viimane katse partitsioneerimise osas oli see, et partitsioneerisin kromosoome ja seejĂ€rel sortisin iga partii. Teoreetiliselt peaks see kiirendama iga pĂ€ringut, kuna soovitud SNP andmed peaksid olema mitme Parquet'i ploki piires antud vahemikus. Kahjuks osutus isegi partitsioneeritud andmete sortimine keeruliseks ĂŒlesandeks. Tulemuseks lĂ€ksin ĂŒle EMR-ile kohandatud klastriga ja kasutasin kaheksat vĂ”imsat instantsi (C5.4xl) ning Sparklyr'i, et luua paindlikum tööprotsess


# 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')
  )


kuid ĂŒlesanne ei olnud siiski lĂ”puks tĂ€idetud. Seadsin asja igat pidi: suurendasin mĂ€lu eraldamist iga pĂ€ringu teostaja jaoks, kasutasin suure mĂ€lu mahuga sĂ”lmi, rakendasin laialehe muutujaid (broadcasting variable), kuid iga kord oli see osaline lahendus ja jĂ€rk-jĂ€rgult hakkasid teostajad ebaĂ”nnestuma, kuni kĂ”ik seisis.

Muudan end leidlikumaks

Mida ma Ôppisin: mÔnikord nÔuavad erilised andmed erilisi lahendusi.

Iga SNP-l on positsioonivÀÀrtus. See on number, mis vastab aluste arvule, mis asuvad selle kromosoomi kĂ”rval. See on hea ja loomulik viis meie andmete korraldamiseks. Alguses tahtsin rĂŒhmitada iga kromosoomi piirkondade jĂ€rgi. NĂ€iteks positsioonid 1–2000, 2001–4000 jne. Kuid probleem on selles, et SNP-d on kromosoomides ebaĂŒhtlaselt jaotatud, seega on grupi suurused oluliselt erinevad.

Parsing 25TB using AWK and R

Tulemusena jĂ”udsin ma positsioonide kategooriate (rank) jagamiseni. Laaditud andmete pĂ”hjal jooksin pĂ€ringu, et saada unikaalsete SNP-de, nende positsioonide ja kromosoomide nimekiri. SeejĂ€rel sorteerisin andmed iga kromosoomi sees ja kogusin SNP-d antud suurusega gruppidesse (bin). Ütleme, 1000 SNP kohta. See andis mulle SNP-de seose grupiga kromosoomis.

LĂ”puks tegin rĂŒhmadesse (bin) 75 SNP, pĂ”hjuse selgitan allpool.

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

Esimene katse Sparkiga

Mida ma Ă”ppisin: ĂŒhendamine Sparkis toimub kiiresti, kuid partitsioneerimine on endiselt kulukas.

Tahtsin laadida selle vĂ€ikese (2,5 miljonit rida) andmeframes Sparkisse, ĂŒhendada selle toorandmetega ja seejĂ€rel partitsioneerida uue lisatud veeru jĂ€rgi. 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')
  )

Olen kasutanud sdf_broadcast(), nii teab Spark, et ta peab andmeframe kĂ”ikidesse sĂ”lmepunktidesse saatma. See on kasulik, kui andmed on vĂ€ikesed ja vajalikud kĂ”igi ĂŒlesannete jaoks. Vastasel juhul pĂŒĂŒab Spark olla nutikas ja jaotab andmeid vastavalt vajadusele, mis vĂ”ib pĂ”hjustada viivitusi.

Ja jĂ€lle ei töötanud mu plaan: ĂŒlesanded töötasid mingil ajal, lĂ”petasid ĂŒhendamise, ja siis, nagu ka kĂ€ivitatud partitsioneerimise tĂ€itjad, hakkasid need ebaĂ”nnestuma.

Lisan AWK

Mida ma Ôppisin: Àrge magage, kui teile Ôpetatakse aluseid. Kindlasti on keegi teie probleemi juba 1980. aastatel lahendanud.

Kuni selle hetkeni oli kÔigi minu ebaÔnnestumiste pÔhjus Sparkis andmete segamine klastris. VÔib-olla on olukorda vÔimalik parandada eelneva töötlemisega. Otsustasin proovida jagada toorandmed kromosoomi veergudeks, nii lootsin pakkuda Sparkile 'eelnevalt partitsioneeritud' andmeid.

Otsisin StackOverflow'ist, kuidas jagada veergude vÀÀrtuste jÀrgi, ja leidsin nii imelise vastuse. AWK abil on vÔimalik jagada tekstifaili veergude vÀÀrtuste jÀrgi, kirjutades skripti ja mitte suunates tulemusi stdout.

Katsena kirjutasin Bash-skripti. Laadisin alla ĂŒhe pakitud TSV faili ja seejĂ€rel dekompressisin selle gzip ja suunasin selle awk.

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

See töötas!

Tuuma tÀitmine

Mida ma Ôppisin: gnu parallel on maagiline asi, kÔik peaksid seda kasutama.

Jagamine toimus ĂŒsna aeglaselt ja kui ma kĂ€ivitasin htop, et kontrollida vĂ”imsuse (ja kallite) EC2 instantsi kasutust, selgus, et kasutan ainult ĂŒhte tuuma ja umbes 200 MB mĂ€lu. Probleemi lahendamiseks, et mitte palju raha kaotada, pidi olema mĂ”eldud, kuidas töö paralleelseks muuta. Õnneks leidsin ma tĂ€iesti hĂ€mmastavast raamatust Data Science at the Command Line Jeroen Janssensilt peatĂŒki, mis oli pĂŒhendatud paralleele töötlemisele. Sealt sain teada gnu parallel, vĂ€ga paindlik meetod mitme lĂ”ime rakendamiseks Unixis.

Parsing 25TB using AWK and R
Kui kĂ€ivitasin jagamise uue protsessi abil, oli kĂ”ik suurepĂ€rane, kuid jĂ€i kitsaskohaks — S3 objektide allalaadimine kettale ei olnud liiga kiire ja ei olnud tĂ€ielikult paralleelne. Selle parandamiseks tegin nii:

  1. Selgitasin vÀlja, et S3 allalaadimise etappi saab otse voos rakendada, tÀielikult vÀlistades vahehoidla kettal. See tÀhendab, et saan vÀltida toorete andmete kirjutamist kettale ja kasutada veel vÀiksemat, seega ka odavamat ladustamist AWS-is.
  2. KÀsuga aws configure set default.s3.max_concurrent_requests 50 suurendas oluliselt lÔime arvu, mida AWS CLI kasutab (vaikimisi on neid 10).
  3. LĂ€ksin ĂŒle vĂ”rgu kaudu optimeeritud EC2 instantsile, mille nimetusel on tĂ€ht n. Avastasin, et arvutusvĂ”ime kadu n-instantside kasutamisel kompenseerib suuresti allalaadimise kiirus.
  4. Muutsin gzip . Tundub, et pigz, see on gzip tööriist, mis suudab teha suurepÀraseid asju algselt paralleelse töötlemise rakendamiseks failide dekompressimisel (see aitas kÔige vÀhem).

# 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

Need sammud on omavahel kombineeritud, et kĂ”ik toimiks vĂ€ga kiiresti. Kiire allalaadimise ja kettale kirjutamise loobumise tĂ”ttu sain nĂŒĂŒd töödelda 5 terabaidist paketti vaid paari tunni jooksul.

Selles tweets pidi olema mainitud 'TSV'. Kahjuks ei olnud.

Uute andmete kasutamine

Mida ma Ôppisin: Spark armastab tihendamata andmeid ja ei salli partitsioonide kombineerimist.

NĂŒĂŒd olid andmed S3-s avatud (loe: jagatud) ja poolkorraldatud formaadis ning ma sain tagasi minna Spark'i juurde. Mind ootas ĂŒllatus: ma ei suutnud jĂ€lle soovitud tulemust saavutada! Oli vĂ€ga keeruline öelda Spark'ile, kuidas andmed on partitsioneeritud. Ja isegi kui ma seda tegin, selgus, et partitsioone on liiga palju (95 tuhat), ja kui ma kasutasin coalesce et vĂ€hendada nende arvu mĂ”istlikesse piiridesse, siis see rikkus minu partitsioneerimise. Olen kindel, et seda saab parandada, kuid paar pĂ€eva otsinguid ei toonud lahendust. LĂ”puks tegin kĂ”ik ĂŒlesanded Spark'is Ă€ra, kuigi see vĂ”ttis aega ja minu jagatud Parquet-failid ei olnud eriti vĂ€ikesed (~200 Kb). Kuid andmed asusid seal, kus pidid.

Parsing 25TB using AWK and R
Liialt vÀikesed ja erinevad, suurepÀrane!

Kohalike Spark'i pÀringute testimine

Mida ma Ă”ppisin: Spark'is on liiga palju kulusid lihtsate ĂŒlesannete tĂ€itmisel.

Andmete ĂŒleslaadimine hĂ€sti mĂ”eldud formaadis vĂ”imaldas mul testida kiirus. Seadsin R-i skripti kohalikku Spark'i serverisse kĂ€ivitamiseks ja seejĂ€rel laadisin Spark'i andmeframe'i mÀÀratud Parquet-gruppide (bin) salvestusest. Proovisem laadida kĂ”ik andmed, kuid ei suutnud Sparklyr'ile partitsioneerimist tuvastada.

sc <- Spark_connect(master = "local")

desired_snp <- 'rs34771739'

# Alusta taimerit
start_time <- Sys.time()

# Laadi soovitud bin Spark'i
intensity_data % 
  Spark_read_Parquet(
    name = 'intensity_data', 
    path = get_snp_location(desired_snp),
    memory = FALSE )

# Alamhulk bin'ist snp ja seejÀrel kogu kohalikusse
test_subset % 
  filter(SNP_Name == desired_snp) %>% 
  collect()

print(Sys.time() - start_time)

Teostamine vÔttis 29,415 sekundit. Palju parem, kuid mitte piisavalt hea massilise testimise jaoks. Pealegi ei saanud ma kiirusetÔusu saavutada mÀlus andmeframe'i vahemÀlustades, kuna iga kord, kui proovisin andmeframe'i mÀlu salvestada, kukkus Spark alati kokku, isegi kui eraldasin rohkem kui 50 GB mÀlu andmestikule, mis kaalus alla 15.

Tagasi AWK'i juurde

Mida ma Ôppisin: assotsiatiivsed massiivid AWK's on vÀga efektiivsed.

MĂ”istsin, et saan saavutada suuremat kiirus. Meenus, et suurepĂ€rases AWK juhendis Bruce Barnett'ilt Ma olin lugenud lahedast funktsioonist, mida nimetatakse „assotsiatiivsed massiivid”. Tegelikult on need vĂ”tme-vÀÀrtuse paarid, mida kuidagi AWK-s teisti nimetatakse, ja seetĂ”ttu ei tulnud ma nendele eriti mĂ”elnud. Roman Cheplyaka tuletas meelde, et mĂ”isted „assotsiatiivsed massiivid” on palju vanemad kui mĂ”isted „vĂ”tme-vÀÀrtuse paar”. Isegi kui te otsite vĂ”tme-vÀÀrtuse paar Google Ngram-st, ei nĂ€e te seda mĂ”istet seal, vaid leiate assotsiatiivsed massiivid! Pealegi seostatakse „vĂ”tme-vÀÀrtuse paare” enamasti andmebaasidega, seega on palju mĂ”istlikum vĂ”rrelda neid hashmapiga. Ma sain aru, et saan neid assotsiatiivseid massiive kasutada oma SNP-de seostamiseks rĂŒhma tabeliga (bin table) ja toorandmete jaoks ilma Sparkita.

Selleks kasutasin AWK skriptis ploki BEGIN. See on koodijupp, mis tÀidetakse enne, kui esimene andmereegel edastatakse skripti peamisse keha.

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"
}

Meeskond while(getline...) laadis kÔik read CSV grupist (bin), mÀÀrates esimese veeru (SNP nimi) vÔtmena assotsiatiivses massiivis bin ja teise vÀÀrtuse (grupp) vÀÀrtusena. SeejÀrel plokis { }, mis tÀidetakse kÔigi peafaile ridade kohta, saadetakse iga rida vÀljastusfaili, millel on ainulaadne nimi, sÔltuvalt tema grupist (bin): ..._bin_"bin[$1]"_....

Muutujad batch_num ja chunk_id vastasid torujuhtme poolt esitatud andmetele, mis vÔimaldas vÀltida vÔidujooksu olukorda, ja igal kÀivitatud parallel, kirjutas oma ainulaadsesse faili.

Kuna kĂ”ik toorandmed olin jaganud kaustadesse kromosoomide jĂ€rgi, mis olid jÀÀnud pĂ€rast minu eelmist AWK eksperimenti, sain nĂŒĂŒd kirjutada teise Bash skripti, et töödelda kromosoomi korraga ja saata S3-sse sĂŒgavamalt partsioneeritud andmeid.

DESIRED_CHR='13'

# Laadi kromosoomi andmed s3-st ja jaga osakesteks
aws s3 ls $DATA_LOC |
awk '{print $4}' |
grep 'chr'$DESIRED_CHR'.csv' |
parallel "echo 'reading {}'; aws s3 cp "$DATA_LOC"{} - | awk -v chr=""$DESIRED_CHR"" -v chunk="{}" -f split_on_chr_bin.awk"

# Koonda kĂ”ik paralleelsed protsessi osad ĂŒhte faili ja laadi RDS-sse, kasutades R-i
ls chunked/ |
cut -d '_' -f 4 |
sort -u |
parallel "echo 'zipping bin {}'; cat chunked/*_bin_{}_*.csv | ./upload_as_rds.R '$S3_DEST'/chr_'$DESIRED_CHR'_bin_{}.rds"
rm chunked/*

Skriptis on kaks osa parallel.

Esimese jaotises loetakse andmed kĂ”igist failidest, mis sisaldavad teavet vajaliku kromosoomi kohta, seejĂ€rel jaotatakse need andmed voogude vahel, mis jaotavad failid vastavatesse gruppidesse (bin). Selleks, et vĂ€ltida seaduse rikkumist, kui mitu voogu kirjutavad ĂŒhte faili, edastab AWK neile ĂŒhenduse failide nimed, kuhu andmeid kirjutada, nĂ€iteks, chr_10_bin_52_batch_2_aa.csv. Tulemusena luuakse kettale palju vĂ€ikseid faile (selleks kasutasin terabaidiseid EBS-mahte).

Teise jaotise torustik parallel lĂ€bib gruppe (bin) ja ĂŒhendab nende eraldi failid ĂŒldiseks CSV-ks cat, seejĂ€rel saadetakse need eksportimiseks.

R-i tÔlkimine?

Mida ma Ôppisin: Saate sellele juurdepÀÀsu stdin ja stdout R-i skriptist, seega saab seda kasutada ka voogus.

Bash-skriptis vĂ”isite sellist rida mĂ€rgata: ...cat chunked/*_bin_{}_*.csv | .\/upload_as_rds.R.... See tĂ”lgib kĂ”ik ĂŒhendatud failid grupist (bin) allolevasse R-skripti. {} on eriline meetod parallel, mis sisestab kĂ”ik selle mÀÀratud voogu saadetavad andmed otse kĂ€susse. Valik {#} pakub unikaalset voogude tĂ€itmise ID-d, ja {%} on tĂ¶Ă¶ĂŒlesande pesade number (kordub, kuid mitte kunagi samaaegselt). KĂ”ikide valikute loendi leiate dokumentatsioonis.

#!/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
)

Kui muutuja file("stdin") edastatakse readr::read_csv, R-skripti tÔlgitud andmed laaditakse raamistikku, mis seejÀrel salvestatakse .rds-failina kasutades aws.s3 salvestatakse otse S3-sse.

RDS on midagi nagu Parquet'i noorem versioon, ilma veergu salvestamise rafineerimiseta.

Bash-skripti lĂ”petamisel sain hulga .rds-faile, mis asusid S3-s, mis vĂ”imaldas mul kasutada efektiivset kokkusurumist ja sisseehitatud tĂŒĂŒpe.

Malest R-i kasutamine töötas vĂ€ga kiiresti, kuigi see oli aeglane. Pole ĂŒllatav, et R-i andmete lugemise ja kirjutamise tĂŒkkide optimeerimine on hea. PĂ€rast keskmise suurusega kromosoomi testimist lĂ”ppes ĂŒlesanne C5n.4xl instansil umbes kahe tunni jooksul.

S3 piirangud

Mida ma Ôppisin: TÀnu nutikale S3 teede rakendusele saab see hallata palju faile.

Ma olin mures, kas S3 suudab hallata suurt hulka edastatud faile. Ma saaksin faili nimed teha mÔistlikeks, kuid kuidas S3 neid otsib?

Parsing 25TB using AWK and R
Kaustad S3-s on vaid ilu pĂ€rast, tegelikult ei huvita sĂŒsteemi sĂŒmbolid /. S3 KKK lehelt.

Tundub, et S3 esindab faili teed lihtsalt vÔtme kaudu omamoodi hajus- vÔi dokumendibaasis. Baki (bucket) vÔib pidada tabeliks ja failid tabeli ridadeks.

Kuna kiirus ja efektiivsus on Amazonis kasumi teenimiseks olulised, ei ole ĂŒllatav, et see 'vĂ”ti-faili-teele' sĂŒsteem on erakordselt optimeeritud. Olen pĂŒĂŒdnud leida tasakaalu: et ei peaks tegema hulgaliselt get-pĂ€ringuid, kuid samas pĂ€ringud toimiksid kiiresti. Selgus, et kĂ”ige parem on teha umbes 20 000 bin-faili. Arvan, et kui jĂ€tkata optimeerimist, vĂ”iks kiirus veelgi tĂ”usta (nĂ€iteks luua spetsiaalne bak, et andmed seal hoida, vĂ€hendades seelĂ€bi otsingutabeli suurust). Kuid edasiste katsetustega ei olnud juba aega ega raha.

Kuidas on lood ristĂŒhilduvusega?

Mida ma Ôppisin: peamine pÔhjus, miks aega kaotatakse, on teie andmete salvestamise meetodi enneaegne optimeerimine.

Selles etapis on vĂ€ga oluline kĂŒsida endalt: 'Miks kasutada omandiĂ”iguse formaate?' PĂ”hjus peitub laadimiskiirusest (gzippitud CSV-failid laadisid 7 korda kauem) ja meie tööprotsesside ĂŒhilduvusest. VĂ”in oma otsuse uuesti lĂ€bi vaadata, kui R suudab hĂ”lpsasti laadida Parquet (vĂ”i Arrow) faile ilma Sparki koormuseta. Meie laboris kasutavad kĂ”ik R-i, ja kui mul on vaja andmeid muuda, siis mul on endiselt allikateks tekstifailid, seega vĂ”in lihtsalt toru uuesti kĂ€ivitada.

Tööde jagamine

Mida ma Ă”ppisin: Ă€rge proovige ĂŒlesandeid kĂ€sitsi optimeerida, las see teeb arvuti.

Olen vooliku ĂŒks kromosoomi tööprotsessi silunud, nĂŒĂŒd tuleb töödelda kĂ”ik ĂŒlejÀÀnud andmed.
Soovisin tĂ”sta mitu EC2 instantsi töötlemiseks, kuid ĂŒhtlasi kartsin saada ÀÀrmiselt tasakaalustamata koormust erinevates töötlemise ĂŒlesannetes (nagu ka Spark kannatas tasakaalustamata partitsioonide all). Lisaks ei olnud mul soov tĂ”sta ĂŒhte instantsi igale kromosoomile, kuna AWS-i kontodele on vaikimisi piirang 10 instantsi.

Siis otsustasin kirjutada R-is skripti töötlemise ĂŒlesannete optimeerimiseks.

Esialgu palusin S3-l arvestada, kui palju ruumi iga kromosoom salvestuses vÔtab.

library(aws.s3)
library(tidyverse)

chr_sizes % 
  mutate(Size = as.numeric(Size)) %>% 
  filter(Size != 0) %>% 
  mutate(
    # Extract chromosome from the file name 
    chr = str_extract(Key, 'chr.{1,4}.csv') %>%
             str_remove_all('chr|.csv')
  ) %>% 
  group_by(chr) %>% 
  summarise(total_size = sum(Size)/1e+9) # Divide to get value in GB



# A 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.
# 
 with 17 more rows

Siis kirjutasin funktsiooni, mis vÔtab kogusumma, segab kromosoomide jÀrjekorra, jagab need gruppidesse num_jobs ja teatab, kui palju eristuvad kÔik töötlemistööd suuruses.

num_jobs <- 7
# Kui suur oleks iga töö, kui jagada ideaalselt?
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)



# A tibble: 1 x 2
     sd data            
             
1  153.

Siis jooksin purrriga tuhat segamist ja valisin parima.

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

Nii sain ma komplekti ĂŒlesandeid, mis olid vĂ€ga sarnased suurusele. SeejĂ€rel jĂ€i vaid oma eelmist Bash-skripti suure ringi alla panna for. Selle optimeerimise kirjutamiseks kulus umbes 10 minutit. Ja see on palju vĂ€hem, kui oleksin kulutanud kĂ€sitsi ĂŒlesannete loomisele, kui need ei oleks tasakaalus. Seega arvan, et sellega eelneva optimeerimisega ei lĂ€inud ma petta.

for DESIRED_CHR in "16" "9" "7" "21" "MT"
do
# Üksiku kromosoomi töötlemise kood
fi

LĂ”pus lisan vĂ€lja lĂŒlitamise kĂ€su:

sudo shutdown -h now


 ja see kĂ”ik Ă”nnestus! AWS CLI abil tĂ”stsin instantsid ja lĂ€bi valiku user_data edastasin neile Bash-skripte nende töötlemise ĂŒlesannete jaoks. Need tĂ€ideti ja suleti automaatselt, nii et ma ei maksnud liigse arvutusvĂ”imsuse eest.

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

Pakendame kokku!

Mida ma Ôppisin: API peab olema lihtne, et tagada lihtsus ja paindlikkus kasutamisel.

LĂ”puks sain ma andmed soovitud kohas ja vormis. JÀÀnud oli protsessi maksimaalselt lihtsustada, et mu kolleegidel oleks kergem. Soovisin luua lihtsa API pĂ€ringute tegemiseks. Kui tulevikus otsustan, et tahan ĂŒle minna .rds Parquet-failide puhul peaks see olema probleem minu jaoks, mitte kolleegide jaoks. SeetĂ”ttu otsustasin teha sisemise R-paketi.

Kogusin ja dokumenteerisin vÀga lihtsa paketi, mis sisaldab vaid mÔnda funktsiooni andmete juurde pÀÀsemiseks, mis on koondatud funktsiooni get_snp. Samuti tegin kolleegidele veebisaidi pkgdown, et nad saaksid hÔlpsasti vaadata nÀiteid ja dokumentatsiooni.

Parsing 25TB using AWK and R

Intelligentne vahemÀlu

Mida ma Ôppisin: kui teie andmed on hÀsti ette valmistatud, on vahemÀllu salvestamine lihtne!

Kuna ĂŒks peamistest töövoogudest rakendas paketi SNP jaoks sama analĂŒĂŒsimudelit, otsustasin kasutada rĂŒhmitamist (binning) oma huvides. SNP-de andmeid edastades on tagastatavale objektile lisatud kogu teave grupist (bin). See tĂ€hendab, et vanad pĂ€ringud vĂ”ivad (teoreetiliselt) uute pĂ€ringute töötlemist kiirendada.

# 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
  }
...

Paketti koostades jooksin palju vĂ”rdlusi, et vĂ”rrelda erinevate meetodite kasutamise kiirus. Soovitan sellega mitte alahinnata, kuna tulemused vĂ”ivad olla ĂŒllatavad. NĂ€iteks dplyr::filter oli mĂ€rgatavalt kiirem ridade jÀÀktingitud filtreerimise kaudu haaramiseks ja ĂŒhe veeru saamine filtreeritud andmekaardist töötas palju kiiremini kui indekseerimise sĂŒntaks.

Pange tĂ€hele, et objekt prev_snp_results sisaldab vĂ”tit snps_in_bin. See on kĂ”igi unikaalsete SNP-de massiiv grupis (bin), mis vĂ”imaldab kiiresti kontrollida, kas eelneva pĂ€ringu andmed on juba olemas. Samuti hĂ”lbustab see kĂ”igi SNP-de rĂŒhma (bin) kaudu tsĂŒklilist lĂ€bimist jĂ€rgmise koodi abil:

# 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 
}

tulemused nÀitasid ainult nelja ebaolulise koodibloki kattuvust, mis olid tingitud POSIX ja ANSI C nÔuetest.

NĂŒĂŒd saame (ja oleme hakanud tĂ”siselt) jooksma mudeleid ja stsenaariume, millele meil varem ei olnud juurdepÀÀsu. Parim asi on see, et mu kolleegidel laboratooriumis ei pea muretsema keerukuste ĂŒle. Neil on lihtsalt töötav funktsioon.

Ja kuigi pakett vabastab nad detailidest, ĂŒritasin teha andmeformaat piisavalt lihtsaks, et nad saaksid sellega toime tulla, kui ma homme Ă€kki kadun 


Kiirus on mĂ€rgatavalt kasvanud. Tavaliselt skaneerime funktsionaalselt olulisi geneetilisi fragmente. Varem ei saanud me seda teha (see oli liiga kallis), kuid nĂŒĂŒd, tĂ€nu grupistruktuurile (bin) ja vahemĂ€lule, kulub ĂŒhe SNP pĂ€ringule keskmiselt vĂ€hem kui 0,1 sekundit ning andmete kasutamine on nii madal, et S3 kulud on tĂŒhised.

KokkuvÔte

Suured andmete töötlustooted on universaalsed. Kui teil on aega, vĂ”ite peaaegu kindlasti kirjutada kiirema lahenduse, rakendades 'nutikat' andmete puhastamist, ladustamist ja ekstraktsioonimeetodeid. LĂ”ppude lĂ”puks kĂ”ik sĂ”ltub kulude ja kasumi analĂŒĂŒsist.

Mida ma Ôppisin:

odavat viisi 25 TB korraga töötlemiseks ei ole;

  • olge ettevaatlik oma Parquet-failide suuruse ja korraldusega;
  • Sparkis peavad partiid olema tasakaalus;
  • kunagi ĐœĐ” ĂŒritage teha 2,5 miljonit partiid;
  • sortimine on endiselt keeruline, nagu ka Spark'i seadistamine;
  • mĂ”nikord vajavad erilised andmed erilisi lahendusi;
  • Sparkis toimib liitmine kiiresti, kuid partitseerimine on endiselt kulukas;
  • Ă€rge magage, kui Ă”petatakse aluseid, keegi on kindlasti teie probleemi juba 1980ndatel lahendanud;
  • — see on maagiline asi, kĂ”ik peaks seda kasutama;
  • gnu parallel Spark armastab kokkusurumata andmeid ja ei armasta partitsioonide ĂŒhendamist;
  • Sparkis on liiga palju kulusid lihtsate ĂŒlesannete lahendamisel;
  • assosiatiivsed massiivid AWK-s on vĂ€ga efektiivsed;
  • neile saab ligi pÀÀseda
  • R-skriptsist, seega saab seda kasutada ka torustikus; stdin ja stdout tĂ€nu nutikale S3 teede rakendusele suudab töödelda palju faile;
  • aegade raiskamise peamine pĂ”hjus on teie ladustamisviisi enneaegne optimeerimine;
  • Ă€rge proovige optimeerida ĂŒlesandeid kĂ€sitsi, laske sellel teha arvuti;
  • API peab olema lihtne, et tagada kasutamise lihtsus ja paindlikkus;
  • kui teie andmed on hĂ€sti ette valmistatud, on vahemĂ€llu salvestamine lihtne!
  • kui teie andmed on hĂ€sti ette valmistatud, on vahemĂ€lu salvestamine lihtne!

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster