
Comment lire cet article: je suis dĂ©solĂ© que le texte soit si long et confus. Pour vous faire gagner du temps, je commence chaque chapitre par une introduction « Ce que j'ai appris », oĂč je rĂ©sume en une ou deux phrases l'essence du chapitre.
« Montre simplement la solution ! » Si vous voulez juste voir oĂč j'en suis, passez au chapitre « Je deviens plus inventive », mais je pense qu'il est plus intĂ©ressant et utile de lire sur les Ă©checs.
RĂ©cemment, j'ai Ă©tĂ© chargĂ© de configurer le processus de traitement d'un grand volume de sĂ©quences d'ADN d'origine (techniquement, il s'agit d'une puce SNP). Je devais rapidement obtenir des donnĂ©es sur une localisation gĂ©nĂ©tique spĂ©cifique (appelĂ©e SNP) pour des modĂ©lisations et d'autres tĂąches. GrĂące Ă R et AWK, j'ai rĂ©ussi Ă nettoyer et organiser les donnĂ©es de maniĂšre naturelle, accĂ©lĂ©rant ainsi considĂ©rablement le traitement des requĂȘtes. Cela nâa pas Ă©tĂ© facile et a nĂ©cessitĂ© de nombreuses itĂ©rations. Cet article vous aidera Ă Ă©viter certaines de mes erreurs et vous montrera ce que j'ai finalement rĂ©ussi Ă accomplir.
Tout d'abord, quelques explications préliminaires.
Données
Notre centre universitaire de traitement de l'information génétique nous a fourni des données au format TSV d'une taille de 25 To. Je les ai reçues sous forme de 5 paquets, compressés en Gzip, chacun contenant environ 240 fichiers de quatre Go. Chaque ligne contenait des données pour un SNP d'une personne. En tout, des données concernant ~2,5 millions de SNP et ~60 000 personnes ont été transmises. En plus des informations sur les SNP, les fichiers comprenaient de nombreuses colonnes de chiffres reflétant diverses caractéristiques, telles que l'intensité de lecture, la fréquence des différents allÚles, etc. Environ 30 colonnes contenaient des valeurs uniques.
L'objectif
Comme dans tout projet de gestion des données, le plus important était de déterminer comment les données seraient utilisées. Dans ce cas nous allons principalement modéliser et établir des flux de travail pour les SNP en fonction des SNP. C'est-à -dire que nous aurons besoin de données uniquement pour un SNP à la fois. Je devais apprendre à extraire le plus simplement, le plus rapidement et le moins cher possible tous les enregistrements relatifs à l'un des 2,5 millions de SNP.
Comment ne pas faire cela
Je vais citer un cliché pertinent :
Je n'ai pas Ă©chouĂ© mille fois, j'ai simplement dĂ©couvert mille façons de ne pas analyser une tonne de donnĂ©es dans un format utile pour les requĂȘtes.
PremiĂšre tentative
Ce que j'ai appris: il n'existe pas de moyen économique de parser 25 To à la fois.
AprĂšs avoir suivi le cours « MĂ©thodes avancĂ©es de traitement des grandes donnĂ©es » Ă l'UniversitĂ© Vanderbilt, j'Ă©tais convaincu que c'Ă©tait facile. Il me faudrait probablement une heure ou deux pour configurer le serveur Hive afin de parcourir toutes les donnĂ©es et de rendre compte des rĂ©sultats. Ătant donnĂ© que nos donnĂ©es sont stockĂ©es sur AWS S3, j'ai utilisĂ© le service , qui permet d'appliquer des requĂȘtes SQL Hive sur les donnĂ©es S3. Pas besoin de configurer ou de lancer un cluster Hive, et vous payez seulement pour les donnĂ©es que vous recherchez.
AprĂšs avoir montrĂ© Ă Athena mes donnĂ©es et leur format, j'ai exĂ©cutĂ© quelques tests avec des requĂȘtes similaires :
select * from intensityData limit 10;Et j'ai rapidement obtenu des rĂ©sultats bien structurĂ©s. PrĂȘt.
Jusqu'Ă ce que nous essayions d'utiliser les donnĂ©es en pratiqueâŠ
On m'a demandĂ© d'extraire toutes les informations sur les SNP afin de tester le modĂšle. J'ai lancĂ© la requĂȘte :
select * from intensityData
where snp = 'rs123456';âŠet j'ai attendu. AprĂšs huit minutes et plus de 4 To de donnĂ©es demandĂ©es, j'ai obtenu un rĂ©sultat. Athena facture en fonction du volume de donnĂ©es trouvĂ©es, Ă 5 $ le tĂ©raoctet. Donc, cette unique requĂȘte m'a coĂ»tĂ© 20 $ et huit minutes d'attente. Pour exĂ©cuter le modĂšle sur toutes les donnĂ©es, il aurait fallu patienter 38 ans et payer 50 millions de dollars. Ăvidemment, cela ne nous convenait pas.
Il fallait utiliser ParquetâŠ
Ce que j'ai appris: faites attention Ă la taille de vos fichiers Parquet et Ă leur organisation.
Tout d'abord, j'ai essayé de résoudre le problÚme en convertissant tous les TSV en . Ils sont pratiques pour travailler avec de grands ensembles de données, car l'information y est stockée sous forme de colonnes : chaque colonne est placée dans son propre segment de mémoire/disque, contrairement aux fichiers texte dans lesquels les lignes contiennent les éléments de chaque colonne. Et quand il faut trouver quelque chose, il suffit de lire la colonne nécessaire. De plus, dans chaque fichier, chaque colonne stocke une plage de valeurs, donc si une valeur recherchée n'est pas présente dans la plage de la colonne, Spark ne perdra pas de temps à balayer tout le fichier.
J'ai lancĂ© une tĂąche simple Pour convertir nos TSV en Parquet, j'ai chargĂ© de nouveaux fichiers dans Athena. Cela a pris environ 5 heures. Mais lorsque j'ai lancĂ© la requĂȘte, son exĂ©cution a pris presque autant de temps et un peu moins d'argent. Le problĂšme est que Spark, essayant d'optimiser la tĂąche, a simplement dĂ©compressĂ© un morceau de TSV et l'a placĂ© dans son propre morceau de Parquet. Et comme chaque morceau Ă©tait assez gros et contenait des enregistrements complets de nombreuses personnes, tous les SNP Ă©taient donc stockĂ©s dans chaque fichier, ce qui obligeait Spark Ă ouvrir tous les fichiers pour extraire les informations nĂ©cessaires.
Il est intĂ©ressant de noter que le type de compression par dĂ©faut (et recommandĂ©) dans Parquet â snappy â n'est pas partageable (splitable). Par consĂ©quent, chaque exĂ©cuteur est restĂ© bloquĂ© sur la tĂąche de dĂ©compression et le chargement de l'ensemble complet de donnĂ©es de 3,5 Go.

Nous examinons le problĂšme
Ce que j'ai appris: le tri est difficile, surtout si les données sont réparties.
J'avais l'impression que je comprenais maintenant l'essentiel du problĂšme. Je devais simplement trier les donnĂ©es par colonne SNP, et non par personnes. Ainsi, plusieurs SNP seraient stockĂ©s dans un seul morceau de donnĂ©es, et alors, la fonction « intelligente » de Parquet « n'ouvre que si la valeur est dans la plage » montrerait toute sa valeur. Malheureusement, trier des milliards de lignes dispersĂ©es sur un cluster s'est avĂ©rĂ© ĂȘtre une tĂąche complexe.
Moi, prenant des cours d'algorithmes à l'université : « Ugh, personne ne se soucie de la complexité computationnelle de tous ces algorithmes de tri »
Moi, essayant de trier sur une colonne dans un 20 To table : « Pourquoi cela prend-il autant de temps ? » difficultés.
â Nick Strayer (@NicholasStrayer)
AWS ne veut certainement pas rendre l'argent pour la raison « Je suis un étudiant distrait ». AprÚs avoir lancé le tri sur Amazon Glue, il a fonctionné pendant 2 jours et a échoué.
Qu'en est-il du partitionnement ?
Ce que j'ai appris: les partitions dans Spark doivent ĂȘtre Ă©quilibrĂ©es.
Puis, j'ai eu l'idée de partitionner les données par chromosomes. Il y en a 23 (et quelques autres, si l'on considÚre l'ADN mitochondrial et les régions non mappées).
Cela permettra de diviser les donnĂ©es en portions plus petites. Si j'ajoute dans la fonction d'exportation de Spark dans le script Glue juste une ligne partition_by = "chr", alors les donnĂ©es devraient ĂȘtre rĂ©parties dans des seaux (buckets).

Le génome est composé de nombreux fragments appelés chromosomes.
Malheureusement, cela n'a pas fonctionnĂ©. Les chromosomes ont des tailles diffĂ©rentes, donc une quantitĂ© d'informations diffĂ©rente. Cela signifie que les tĂąches envoyĂ©es par Spark aux workers n'Ă©taient pas Ă©quilibrĂ©es et s'exĂ©cutaient lentement, car certains nĆuds terminaient plus tĂŽt et restaient inactifs. Cependant, les tĂąches ont Ă©tĂ© exĂ©cutĂ©es. Mais lors de la demande d'un SNP, le dĂ©sĂ©quilibre a de nouveau causĂ© des problĂšmes. Le coĂ»t de traitement des SNP dans des chromosomes plus grands (c'est-Ă -dire ceux d'oĂč nous voulons obtenir des donnĂ©es) n'a diminuĂ© que d'environ 10 fois. Beaucoup, mais pas assez.
Et si l'on les divisait en partitions encore plus petites ?
Ce que j'ai appris: ne tentez jamais de créer 2,5 millions de partitions.
J'ai dĂ©cidĂ© de faire les choses en grand et de partitionner chaque SNP. Cela a garanti une taille de partitions uniforme. C'ĂTAIT UNE MAUVAISE IDĂE. J'ai utilisĂ© Glue et ajoutĂ© une simple ligne partition_by = 'snp'. Le travail a Ă©tĂ© lancĂ© et a commencĂ© Ă s'exĂ©cuter. Un jour plus tard, j'ai vĂ©rifiĂ© et j'ai vu qu'il n'y avait toujours rien Ă©crit dans S3, donc j'ai arrĂȘtĂ© la tĂąche. Il semble que Glue Ă©crivait des fichiers intermĂ©diaires dans un emplacement cachĂ© dans S3, avec beaucoup de fichiers, peut-ĂȘtre plusieurs millions. En consĂ©quence, mon erreur m'a coĂ»tĂ© plus de mille dollars et n'a pas rĂ©joui mon mentor.
Partitionnement + tri
Ce que j'ai appris: trier reste difficile, tout comme configurer Spark.
La derniĂšre tentative de partitionnement consistait Ă partitionner les chromosomes, puis Ă trier chaque partition. En thĂ©orie, cela permettrait d'accĂ©lĂ©rer chaque requĂȘte, car les donnĂ©es de SNP souhaitĂ©es devraient se trouver dans quelques morceaux Parquet Ă l'intĂ©rieur d'une plage donnĂ©e. Malheureusement, le tri des donnĂ©es mĂȘme partitionnĂ©es s'est rĂ©vĂ©lĂ© ĂȘtre une tĂąche difficile. En consĂ©quence, je suis passĂ© Ă EMR pour un cluster personnalisĂ© et j'ai utilisĂ© huit instances puissantes (C5.4xl) et Sparklyr pour crĂ©er un flux de travail plus flexibleâŠ
# 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')
)⊠cependant, la tĂąche n'a toujours pas Ă©tĂ© exĂ©cutĂ©e. J'ai essayĂ© diverses configurations : augmenter l'attribution de mĂ©moire pour chaque exĂ©cuteur de requĂȘtes, utiliser des nĆuds avec plus de mĂ©moire, appliquer des variables de diffusion (broadcasting variables), mais Ă chaque fois cela s'est rĂ©vĂ©lĂ© ĂȘtre des solutions temporaires, et progressivement les exĂ©cuteurs ont commencĂ© Ă Ă©chouer jusqu'Ă ce que tout s'arrĂȘte.
Mise à jour : ça commence.
â Nick Strayer (@NicholasStrayer)
Je deviens plus inventif
Ce que j'ai appris: parfois, des données spécifiques nécessitent des solutions spécifiques.
Chaque SNP a une valeur de position. Ce nombre correspond à la quantité de bases le long de son chromosome. C'est une bonne et naturelle façon d'organiser nos données. Au départ, je voulais partitionner par zones de chaque chromosome. Par exemple, positions 1 - 2000, 2001 - 4000, etc. Mais le problÚme est que les SNPs sont répartis de maniÚre inégale sur les chromosomes, donc la taille des groupes va varier considérablement.

En consĂ©quence, j'en suis venu Ă une classification par catĂ©gories (rang) des positions. Sur les donnĂ©es dĂ©jĂ chargĂ©es, j'ai exĂ©cutĂ© une requĂȘte pour obtenir une liste unique de SNPs, leurs positions et chromosomes. Puis, j'ai triĂ© les donnĂ©es Ă l'intĂ©rieur de chaque chromosome et regroupĂ© les SNPs en groupes (bin) de taille donnĂ©e. Disons, par exemple, 1000 SNPs. Cela m'a donnĂ© une relation SNP avec un groupe dans le chromosome.
Finalement, j'ai créé des groupes (bin) de 75 SNPs, la raison sera expliquée ci-dessous.
snp_to_bin %
group_by(chr) %>%
arrange(position) %>%
mutate(
rank = 1:n()
bin = floor(rank/snps_per_bin)
) %>%
ungroup()PremiĂšre tentative avec Spark
Ce que j'ai appris: la jointure dans Spark fonctionne rapidement, mais la partition reste coûteuse.
Je voulais lire ce petit (2,5 millions de lignes) dataframe dans Spark, le joindre avec des données brutes, puis le partitionner selon la nouvelle colonne ajoutée. 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')
) J'ai utilisĂ© sdf_broadcast(), de cette maniĂšre Spark sait qu'il doit envoyer le dataframe Ă tous les nĆuds. Cela est utile si les donnĂ©es sont de petite taille et nĂ©cessaires pour toutes les tĂąches. Sinon, Spark essaie d'ĂȘtre intelligent et rĂ©partit les donnĂ©es au besoin, ce qui peut causer des ralentissements.
Et encore une fois, mon idée n'a pas fonctionné : les tùches ont fonctionné un moment, ont terminé la jointure, puis, tout comme les exécuteurs lancés par la partition, ont commencé à échouer.
J'ajoute AWK
Ce que j'ai appris: ne dormez pas quand on vous enseigne les bases. Ăvidemment, quelqu'un a dĂ©jĂ rĂ©solu votre problĂšme dans les annĂ©es 1980.
Jusqu'Ă prĂ©sent, la raison de tous mes Ă©checs avec Spark Ă©tait la dĂ©sorganisation des donnĂ©es dans le cluster. Peut-ĂȘtre que la situation peut ĂȘtre amĂ©liorĂ©e par un prĂ©traitement. J'ai dĂ©cidĂ© d'essayer de diviser les donnĂ©es textuelles brutes en colonnes de chromosomes, espĂ©rant ainsi fournir Ă Spark des donnĂ©es « prĂ©-partitionnĂ©es ».
J'ai cherché sur StackOverflow comment diviser par les valeurs des colonnes et j'ai trouvé Avec AWK, vous pouvez diviser un fichier texte en fonction des valeurs des colonnes en exécutant un script, plutÎt qu'en envoyant les résultats à stdout.
Pour faire un essai, j'ai écrit un script Bash. J'ai téléchargé l'un des fichiers TSV compressés, puis je l'ai décompressé avec gzip et envoyé à awk.
gzip -dc path/to/chunk/file.gz |
awk -F 't'
'{print $1",..."$30" > "chunked/"$chr"_chr"$15".csv"}'Ăa a fonctionnĂ© !
Remplissage des cĆurs
Ce que j'ai appris: gnu parallel â c'est une chose magique, tout le monde devrait l'utiliser.
Le division se faisait assez lentement, et lorsque j'ai lancĂ© htop, pour vĂ©rifier l'utilisation d'une instance EC2 puissante (et coĂ»teuse), il s'est avĂ©rĂ© que j'utilisais seulement un cĆur et environ 200 Mo de mĂ©moire. Pour rĂ©soudre le problĂšme sans perdre une fortune, il fallait penser Ă comment parallĂ©liser le travail. Heureusement, dans le livre absolument incroyable de Jeroen Janssens, j'ai trouvĂ© un chapitre dĂ©diĂ© Ă la parallĂ©lisation. J'y ai appris Ă propos de gnu parallel, une mĂ©thode trĂšs flexible d'implĂ©mentation de la multithreading sous Unix.

Lorsque j'ai lancé la division avec le nouveau processus, tout était parfait, mais il restait un goulet d'étranglement: le téléchargement des objets S3 sur le disque n'était pas trÚs rapide et pas entiÚrement parallélisé. Pour corriger cela, j'ai fait ce qui suit :
- J'ai découvert qu'il était possible d'implémenter l'étape de téléchargement S3 directement dans le pipeline, éliminant complÚtement le stockage intermédiaire sur disque. Cela signifie que je peux éviter d'écrire des données brutes sur le disque et utiliser un stockage encore plus petit, et donc moins cher, sur AWS.
- Avec la commande
aws configure set default.s3.max_concurrent_requests 50a fortement augmenté le nombre de threads utilisés par l'AWS CLI (par défaut 10). - Je suis passé à une instance EC2 optimisée pour la vitesse du réseau, avec la lettre n dans le nom. J'ai découvert que la perte de puissance de calcul lors de l'utilisation d'instances n était largement compensée par l'augmentation de la vitesse de téléchargement. Pour la plupart des tùches, j'ai utilisé c5n.4xl.
- J'ai changé
gzipsur , c'est un outil gzip qui peut faire des choses intéressantes pour paralléliser une tùche non parallélisée de décompression de fichiers (cela a aidé le moins).
# 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/*
doneCes étapes sont combinées les unes avec les autres pour que tout fonctionne trÚs rapidement. Grùce à l'augmentation de la vitesse de téléchargement et à l'élimination de l'écriture sur disque, j'ai maintenant pu traiter un lot de 5 To en quelques heures.
Il n'y a rien de plus doux que de voir tous les cĆurs que vous payez sur AWS ĂȘtre utilisĂ©s. GrĂące Ă gnu-parallel, je peux dĂ©compresser et diviser un fichier csv de 19 Go aussi rapidement que je peux le tĂ©lĂ©charger. Je n'ai mĂȘme pas pu faire fonctionner Spark.
â Nick Strayer (@NicholasStrayer)
Ce tweet devait mentionner 'TSV'. Malheureusement.
Réutilisation des données analysées
Ce que j'ai appris: Spark préfÚre les données non compressées et n'aime pas combiner les partitions.
Les donnĂ©es Ă©taient maintenant dans S3 au format non compressĂ© (c'est-Ă -dire dĂ©limitĂ©) et semi-ordonnĂ©, et je pouvais revenir Ă Spark. J'ai eu une surprise : je n'ai pas rĂ©ussi Ă obtenir le rĂ©sultat souhaitĂ© ! Il Ă©tait trĂšs difficile de dire avec prĂ©cision Ă Spark comment les donnĂ©es Ă©taient partitionnĂ©es. MĂȘme lorsque je l'ai fait, il s'est avĂ©rĂ© que le nombre de partitions Ă©tait trop Ă©levĂ© (95 000), et lorsque j'ai utilisĂ© coalesce pour rĂ©duire leur nombre Ă des limites raisonnables, cela a nuit Ă ma partition. Je suis sĂ»r que cela peut ĂȘtre corrigĂ©, mais aprĂšs quelques jours de recherche, je n'ai pas trouvĂ© de solution. Finalement, j'ai rĂ©ussi Ă terminer toutes les tĂąches dans Spark, mĂȘme si cela a pris un certain temps, et mes fichiers Parquet divisĂ©s n'Ă©taient pas trĂšs petits (~200 Ko). Cependant, les donnĂ©es Ă©taient lĂ oĂč elles devaient ĂȘtre.

Trop petites et inégales, c'est merveilleux !
Test des requĂȘtes Spark locales
Ce que j'ai appris: il y a trop de surcharge dans Spark pour résoudre des tùches simples.
En chargeant les données dans un format réfléchi, j'ai pu tester la vitesse. J'ai configuré un script en R pour exécuter un serveur Spark local, puis j'ai chargé un dataframe Spark depuis le stockage Parquet spécifié (bin). J'ai essayé de charger toutes les données, mais je n'ai pas pu faire reconnaßtre la partition à Sparklyr.
sc <- Spark_connect(master = "local")
desired_snp <- 'rs34771739'
# Démarrer un chronomÚtre
start_time <- Sys.time()
# Charger le bin désiré dans Spark
intensity_data %
Spark_read_Parquet(
name = 'intensity_data',
path = get_snp_location(desired_snp),
memory = FALSE )
# Sous-ensemble bin pour snp puis collecter localement
test_subset %
filter(SNP_Name == desired_snp) %>%
collect()
print(Sys.time() - start_time)L'exĂ©cution a pris 29,415 secondes. Beaucoup mieux, mais pas assez bon pour des tests massifs. De plus, je ne pouvais pas accĂ©lĂ©rer le travail grĂące au cache, car lorsque j'essayais de mettre en cache le dataframe en mĂ©moire, Spark plantait toujours, mĂȘme lorsque j'ai allouĂ© plus de 50 Go de mĂ©moire pour un dataset qui pesait moins de 15.
Retour Ă AWK
Ce que j'ai appris: les tableaux associatifs dans AWK sont trĂšs efficaces.
Je savais que je pouvais obtenir une vitesse plus Ă©levĂ©e. Je me suis rappelĂ© que dans le magnifique J'ai lu Ă propos d'une super fonctionnalitĂ© appelĂ©e «». En fait, ce sont des paires clĂ©-valeur, qui, pour une raison quelconque, ont un autre nom dans AWK, et c'est pourquoi je n'y pense pas vraiment. m'a rappelĂ© que le terme «tableaux associatifs» est beaucoup plus ancien que le terme «paire clĂ©-valeur». MĂȘme si vous , vous ne trouverez pas ce terme lĂ -bas, mais vous trouverez des tableaux associatifs ! De plus, «paire clĂ©-valeur» est le plus souvent associĂ© aux bases de donnĂ©es, il est donc beaucoup plus logique de comparer avec un hashmap. J'ai compris que je pouvais utiliser ces tableaux associatifs pour lier mes SNP Ă la table de groupes (bin table) et aux donnĂ©es brutes sans utiliser Spark.
Pour cela, dans le script AWK, j'ai utilisé le bloc BEGIN. C'est un fragment de code qui s'exécute avant que la premiÚre ligne de données ne soit transmise au corps principal du script.
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"
} Commande while(getline...) a chargé toutes les lignes du CSV de groupe (bin), a défini la premiÚre colonne (nom du SNP) comme clé pour le tableau associatif bin et la deuxiÚme valeur (groupe) comme valeur. Ensuite, dans le bloc { }, qui s'exécute pour toutes les lignes du fichier principal, chaque ligne est envoyée dans un fichier de sortie, qui reçoit un nom unique en fonction de son groupe (bin) : ..._bin_"bin[$1]"_....
Variables batch_num et chunk_id correspondaient aux données fournies par le pipeline, ce qui a permis d'éviter une condition de course, et chaque thread d'exécution, lancé en parallÚle, écrivait dans son propre fichier unique.
Puisque j'avais réparti toutes les données brutes dans des dossiers par chromosomes, restés aprÚs ma précédente expérience avec AWK, je pouvais maintenant écrire un autre script Bash pour traiter un chromosome à la fois et fournir des données plus profondément partitionnées dans S3.
DESIRED_CHR='13'
# Téléchargez les données de chromosome depuis s3 et divisez-les en bins
aws s3 ls $DATA_LOC |
awk '{print $4}' |
grep 'chr'$DESIRED_CHR'.csv' |
parallel "echo 'lecture {}'; aws s3 cp "$DATA_LOC"{} - | awk -v chr=""$DESIRED_CHR"" -v chunk="{}" -f split_on_chr_bin.awk"
# Combinez tous les morceaux de processus parallÚles en fichiers uniques et téléchargez-les sur rds en utilisant R
ls chunked/ |
cut -d '_' -f 4 |
sort -u |
parallel "echo 'compression bin {}'; cat chunked/*_bin_{}_*.csv | ./upload_as_rds.R '$S3_DEST'/chr_'$DESIRED_CHR'_bin_{}.rds"
rm chunked/* Le script a deux sections en parallĂšle.
Dans la premiĂšre section, les donnĂ©es sont lues Ă partir de tous les fichiers contenant des informations sur le chromosome requis, puis ces donnĂ©es sont rĂ©parties sur des threads, qui rĂ©partissent les fichiers dans les groupes correspondants (bin). Pour Ă©viter une condition de course, lorsque plusieurs threads Ă©crivent dans un mĂȘme fichier, AWK fournit des noms de fichiers pour Ă©crire des donnĂ©es Ă diffĂ©rents emplacements, par exemple, chr_10_bin_52_batch_2_aa.csv. En consĂ©quence, un grand nombre de petits fichiers sont créés sur le disque (j'ai utilisĂ© des volumes EBS d'un tĂ©raoctet).
Le pipeline de la deuxiĂšme section en parallĂšle parcourt les groupes (bin) et fusionne leurs fichiers individuels en un CSV commun avec cat, puis les envoie Ă l'exportation.
Transcription dans R ?
Ce que j'ai appris: on peut y accéder depuis stdin et stdout un script R, ce qui signifie qu'on peut l'utiliser dans le pipeline.
Dans le script Bash, vous avez peut-ĂȘtre remarquĂ© cette ligne : ...cat chunked/*_bin_{}_*.csv | ./upload_as_rds.R.... Elle transcrit tous les fichiers concatĂ©nĂ©s du groupe (bin) dans le script R ci-dessous. {} est une mĂ©thode spĂ©ciale en parallĂšle, qui insĂšre toutes les donnĂ©es envoyĂ©es dans le flux spĂ©cifiĂ© directement dans la commande elle-mĂȘme. L'option {#} fournit un ID unique de flux d'exĂ©cution, et {%} reprĂ©sente le numĂ©ro de slot de tĂąche (se rĂ©pĂšte, mais jamais en mĂȘme temps). Une liste de toutes les options peut ĂȘtre trouvĂ©e dans
#!/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
) Lorsque la variable file("stdin") est transmise à readr::read_csv, les données transmises au script R sont chargées dans un cadre qui est ensuite enregistré en tant que .rds-fichier à l'aide de aws.s3 , enregistré directement dans S3.
RDS est comme une version simplifiée de Parquet, sans les raffinements d'un entrepÎt en colonnes.
AprÚs la fin du script Bash, j'ai obtenu une série .rds-de fichiers situés dans S3, ce qui m'a permis d'utiliser une compression efficace et des types incorporés.
Malgré l'utilisation de l'R, tout fonctionnait trÚs rapidement. Il n'est pas surprenant que les fragments en R, responsables de la lecture et de l'écriture des données, soient bien optimisés. AprÚs des tests sur un chromosome de taille moyenne, la tùche s'est réalisée sur une instance C5n.4xl en environ deux heures.
Limitations de S3
Ce que j'ai appris: grùce à une implémentation intelligente des chemins, S3 peut gérer un grand nombre de fichiers.
Je m'inquiétais de savoir si S3 pourrait traiter le grand nombre de fichiers qui lui étaient transmis. Je pouvais donner des noms de fichiers significatifs, mais comment S3 les rechercherait-elle ?

Les dossiers dans S3 ne sont que pour l'esthétique, en réalité, le systÚme ne se soucie pas du symbole /.
Il semble que S3 reprĂ©sente le chemin vers un fichier spĂ©cifique sous la forme d'une clĂ© dans une sorte de table de hachage ou base de donnĂ©es basĂ©e sur des documents. Un bucket peut ĂȘtre considĂ©rĂ© comme une table, et les fichiers comme des enregistrements dans cette table.
Puisque la rapiditĂ© et l'efficacitĂ© sont cruciales pour gĂ©nĂ©rer des bĂ©nĂ©fices sur Amazon, il n'est pas surprenant que ce systĂšme de « clĂ©-en-tant-que-chemin-pour-le-fichier » soit extrĂȘmement optimisĂ©. J'ai essayĂ© de trouver un Ă©quilibre : il ne fallait pas effectuer de nombreuses demandes GET, tout en veillant Ă ce que les requĂȘtes soient traitĂ©es rapidement. Il s'est avĂ©rĂ© qu'il Ă©tait prĂ©fĂ©rable de faire environ 20 000 fichiers binaires. Je pense que si l'on continue Ă optimiser, on pourrait augmenter la vitesse (par exemple, en crĂ©ant un bucket spĂ©cial uniquement pour les donnĂ©es, rĂ©duisant ainsi la taille de la table de recherche). Mais il n'y avait dĂ©jĂ plus de temps ni d'argent pour d'autres expĂ©riences.
Qu'en est-il de la compatibilité croisée ?
Ce que j'ai appris : la principale raison de perte de temps est l'optimisation prématurée de votre méthode de stockage.
à ce moment-là , il est trÚs important de se poser la question : « Pourquoi utiliser un format de fichier propriétaire ? » La raison réside dans la rapidité de chargement (les fichiers CSV compressés en gzip prenaient 7 fois plus de temps à se charger) et la compatibilité avec nos flux de travail. Je peux revoir ma décision si R peut facilement charger des fichiers Parquet (ou Arrow) sans la charge de Spark. Dans notre laboratoire, tout le monde utilise R, et si je dois convertir les données dans un autre format, j'ai toujours les données textuelles d'origine, donc je peux simplement relancer le pipeline.
Division du travail
Ce que j'ai appris: ne tentez pas d'optimiser les tĂąches manuellement, laissez cela Ă l'ordinateur.
J'ai validé le workflow sur un chromosome, maintenant il faut traiter toutes les autres données.
Je voulais lancer plusieurs instances EC2 pour la transformation, mais en mĂȘme temps, j'avais peur d'obtenir une charge fortement dĂ©sĂ©quilibrĂ©e dans les diffĂ©rentes tĂąches de traitement (tout comme Spark souffrait de partitions dĂ©sĂ©quilibrĂ©es). De plus, je n'Ă©tais pas ravi de lancer une instance pour chaque chromosome, car il y a une limite par dĂ©faut de 10 instances pour les comptes AWS.
J'ai donc décidé d'écrire un script en R pour optimiser les tùches de traitement.
D'abord, j'ai demandé à S3 de calculer combien d'espace de stockage chaque chromosome occupait.
library(aws.s3)
library(tidyverse)
chr_sizes %
mutate(Size = as.numeric(Size)) %>%
filter(Size != 0) %>%
mutate(
# Extraire le chromosome du nom de fichier
chr = str_extract(Key, 'chr.{1,4}.csv') %>%
str_remove_all('chr|.csv')
) %>%
group_by(chr) %>%
summarise(total_size = sum(Size)/1e+9) # Diviser pour obtenir la valeur en Go
# Un 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.
# ⊠avec 17 autres lignes Ensuite, j'ai écrit une fonction qui prend la taille totale, mélange l'ordre des chromosomes et les divise en groupes. num_jobs et rapporte à quel point les tailles de tous les travaux de traitement varient.
num_jobs <- 7
# Quelle serait la taille de chaque travail si elle était parfaitement divisée ?
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)
# Un tibble: 1 x 2
sd data
1 153.Puis j'ai exécuté avec purrr mille mélanges et sélectionné le meilleur.
1:1000 %>%
map_df(shuffle_job) %>%
filter(sd == min(sd)) %>%
pull(data) %>%
pluck(1) Ainsi, j'ai obtenu un ensemble de travaux trÚs similaires en taille. Il ne restait plus qu'à envelopper mon script Bash précédent dans une grande boucle. for. L'écriture de cette optimisation a pris environ 10 minutes. Et c'est beaucoup moins que ce que j'aurais passé à créer manuellement des travaux en cas de déséquilibre. C'est pourquoi je pense que cette optimisation préliminaire était une bonne décision.
for DESIRED_CHR in "16" "9" "7" "21" "MT"
do
# Code pour traiter un seul chromosome
fiĂ la fin, j'ajoute la commande d'arrĂȘt :
sudo shutdown -h now ⊠et tout a bien fonctionné ! Avec AWS CLI, j'ai lancé des instances et via l'option user_data j'ai transmis des scripts Bash pour leurs travaux de traitement. Ils s'exécutaient et s'éteignaient automatiquement, donc je ne payais pas pour une puissance de calcul excédentaire.
aws ec2 run-instances ...
--tag-specifications "ResourceType=instance,Tags=[{Key=Name,Value=<>}]"
--user-data file://<>On emballe !
Ce que j'ai appris: L'API doit ĂȘtre simple pour la simplicitĂ© et la flexibilitĂ© d'utilisation.
Enfin, j'ai obtenu les donnĂ©es au bon endroit et dans le bon format. Il ne restait plus qu'Ă simplifier au maximum le processus d'utilisation des donnĂ©es pour faciliter la tĂąche de mes collĂšgues. Je voulais crĂ©er une API simple pour faire des requĂȘtes. Si je dĂ©cide un jour de passer Ă .rds Pour les fichiers Parquet, cela devrait ĂȘtre un problĂšme pour moi et non pour mes collĂšgues. Pour cela, j'ai dĂ©cidĂ© de crĂ©er un package R interne.
J'ai assemblé et documenté un package trÚs simple contenant seulement quelques fonctions pour accéder aux données, regroupées autour de la fonction get_snp. J'ai également créé pour mes collÚgues un site , afin qu'ils puissent facilement consulter les exemples et la documentation.

Mise en cache intelligente
Ce que j'ai appris: si vos données sont bien préparées, il sera facile de mettre en cache !
Comme un des principaux flux de travail appliquait au package SNP le mĂȘme modĂšle d'analyse, j'ai dĂ©cidĂ© d'utiliser le regroupement Ă mon avantage. Lors de la transmission des donnĂ©es par SNP, toutes les informations du groupe sont attachĂ©es Ă l'objet retournĂ©. Cela signifie que les anciennes requĂȘtes peuvent (en thĂ©orie) accĂ©lĂ©rer le traitement des nouvelles requĂȘtes.
# 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
}
... Lors de la crĂ©ation du package, j'ai exĂ©cutĂ© de nombreux benchmarks pour comparer la vitesse en utilisant diffĂ©rentes mĂ©thodes. Je recommande de ne pas nĂ©gliger cela, car parfois les rĂ©sultats peuvent ĂȘtre surprenants. Par exemple, dplyr::filter s'est avĂ©rĂ© beaucoup plus rapide que la capture des lignes par filtrage basĂ© sur l'indexation, et obtenir une colonne unique Ă partir d'un cadre de donnĂ©es filtrĂ© fonctionnait beaucoup plus rapidement qu'en utilisant une syntaxe d'indexation.
Notez que l'objet prev_snp_results contient la clĂ© snps_in_bin. Il s'agit d'un tableau de tous les SNP uniques dans le groupe, permettant de vĂ©rifier rapidement si des donnĂ©es de requĂȘtes prĂ©cĂ©dentes existent dĂ©jĂ . Cela simplifie Ă©galement le passage cyclique Ă travers tous les SNP du groupe avec ce code :
# 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
}Résultats
Nous pouvons maintenant (et avons commencé sérieusement) à exécuter des modÚles et des scénarios qui étaient auparavant inaccessibles. Le meilleur, c'est que mes collÚgues de laboratoire n'ont pas à se soucier des complexités. Ils ont simplement une fonction qui fonctionne.
Et bien que le package les libĂšre des dĂ©tails, j'ai essayĂ© de rendre le format des donnĂ©es suffisamment simple pour qu'ils puissent comprendre, au cas oĂč je disparaĂźtrais soudainement demainâŠ
La vitesse a considĂ©rablement augmentĂ©. En gĂ©nĂ©ral, nous analysons des fragments fonctionnellement significatifs du gĂ©nome. Auparavant, nous ne pouvions pas le faire (c'Ă©tait trop coĂ»teux), mais maintenant, grĂące Ă la structure de groupe et Ă la mise en cache, une requĂȘte pour un SNP prend en moyenne moins de 0,1 seconde, et l'utilisation des donnĂ©es est si faible que les coĂ»ts en deviennent nĂ©gligeables.
Récemment, j'ai été chargé de gérer plus de 25 To de données de génotypage brutes pour mon laboratoire. Au début, utiliser Spark prenait 8 minutes et coûtait 20 $ pour interroger un SNP. AprÚs avoir utilisé AWK + pour traiter, cela prend désormais moins d'un dixiÚme de seconde et coûte 0,00001 $. Mon gain personnel win.
â Nick Strayer (@NicholasStrayer)
Conclusion
Cet article n'est pas un guide. La solution est devenue individuelle et est presque certainement non optimale. C'est plutĂŽt un rĂ©cit d'un parcours. Je veux que d'autres comprennent que de telles solutions ne naissent pas entiĂšrement formĂ©es dans l'esprit, mais sont le rĂ©sultat d'essais et d'erreurs. De plus, si vous cherchez un spĂ©cialiste en analyse de donnĂ©es, gardez Ă l'esprit que pour utiliser efficacement ces outils, il faut de l'expĂ©rience, et l'expĂ©rience coĂ»te de l'argent. Je suis heureux d'avoir eu les moyens de paiement, mais beaucoup d'autres, qui pourraient faire le mĂȘme travail mieux que moi, n'auront jamais cette opportunitĂ© en raison d'un manque d'argent, mĂȘme pour une tentative.
Les outils pour les grandes données sont universels. Si vous avez le temps, vous pourrez certainement écrire une solution plus rapide en appliquant un 'nettoyage' intelligent des données, un stockage et des méthodes d'extraction. En fin de compte, tout se résume à une analyse des coûts et des avantages.
Ce que j'ai appris :
- il n'existe pas de moyen peu coûteux de parser 25 To à la fois ;
- faites attention Ă la taille de vos fichiers Parquet et Ă leur organisation ;
- les partitions dans Spark doivent ĂȘtre Ă©quilibrĂ©es ;
- ne jamais essayer de faire 2,5 millions de partitions ;
- le tri reste difficile, tout comme la configuration de Spark ;
- parfois, des données spécifiques nécessitent des solutions spécifiques ;
- la fusion dans Spark fonctionne rapidement, mais le partitionnement reste coûteux ;
- ne dormez pas pendant qu'on vous enseigne les bases, quelqu'un a probablement déjà résolu votre problÚme dans les années 1980 ;
gnu parallelâ c'est une chose magique, tout le monde devrait l'utiliser ;- Spark aime les donnĂ©es non compressĂ©es et n'aime pas combiner les partitions ;
- il y a trop de surcharge dans Spark pour résoudre des tùches simples ;
- les tableaux associatifs dans AWK sont trĂšs efficaces ;
- vous pouvez appeler
stdinetstdoutdepuis un script R, ce qui signifie que vous pouvez l'utiliser dans votre pipeline ; - grĂące Ă une mise en Ćuvre intelligente des chemins S3 pouvant traiter de nombreux fichiers ;
- la principale raison de la perte de temps est l'optimisation prématurée de votre méthode de stockage ;
- n'essayez pas d'optimiser les tĂąches manuellement, laissez l'ordinateur le faire ;
- l'API doit ĂȘtre simple pour facilitĂ© et flexibilitĂ© d'utilisation ;
- si vos données sont bien préparées, le caching sera facile !
Source : habr.com
