
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 , 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 . 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 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.

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.
Mina vĂ”tmas algoritmide kursust kolledĆŸis: «Kehv, kedagi ei huvita nende kĂ”igi sortimisalgoritmide arvutustĂŒĂŒbid»
Mina ĂŒritamas sortida 20TB-s veerus tabel: «Miks see nii kaua aega vĂ”tab?» vaevad.
â Nick Strayer (@NicholasStrayer)
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).

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.
Uuendus: see algab.
â Nick Strayer (@NicholasStrayer)
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.

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

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:
- 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.
- KĂ€suga
aws configure set default.s3.max_concurrent_requests 50suurendas oluliselt lĂ”ime arvu, mida AWS CLI kasutab (vaikimisi on neid 10). - 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.
- Muutsin
gzip. Tundub, et , 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/*
doneNeed 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.
Mitte pole midagi magusamat kui nÀha, et kÔik teie AWS-is tasutud tuumad on kasutuses. TÀnu gnu-parallelile saan ma 19gig CSV faili sama kiiresti lahti pakkida ja jagada kui seda alla laadida. Ma ei suutnud isegi Spark'i tööle panna.
â Nick Strayer (@NicholasStrayer)
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.

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 Ma olin lugenud lahedast funktsioonist, mida nimetatakse ââ. Tegelikult on need vĂ”tme-vÀÀrtuse paarid, mida kuidagi AWK-s teisti nimetatakse, ja seetĂ”ttu ei tulnud ma nendele eriti mĂ”elnud. tuletas meelde, et mĂ”isted âassotsiatiivsed massiividâ on palju vanemad kui mĂ”isted âvĂ”tme-vÀÀrtuse paarâ. Isegi kui te , 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
#!/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?

Kaustad S3-s on vaid ilu pĂ€rast, tegelikult ei huvita sĂŒsteemi sĂŒmbolid /.
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
fiLĂ”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 , et nad saaksid hÔlpsasti vaadata nÀiteid ja dokumentatsiooni.

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.
Hiljuti sain vastutada ĂŒle 25 TB toore genotĂŒĂŒbiandmete haldamise eest oma laboris. Kui alustasin, nĂ”udis spark SNP pĂ€ringuks 8 minutit ja maksis $20. PĂ€rast AWK + töötlemise kasutamist vĂ”tab see nĂŒĂŒd vĂ€hem kui 1/10 sekundi ja maksab $0,00001. Minu isiklik vĂ”it. pic.twitter.com/ANOXVGrmkk
â Nick Strayer (@NicholasStrayer)
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 parallelSpark 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;
stdinjastdouttÀ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
