Le défi du milliard de lignes
(morling.dev)- Le One Billion Row Challenge (1BRC), organisé pendant tout le mois de janvier 2024, était un défi de performance consistant à traiter un fichier texte d’un milliard de lignes pour voir jusqu’où Java pouvait être poussé en vitesse
- L’entrée est un simple texte au format
station;temperature, mais il faut calculer les températures minimale, moyenne et maximale pour chaque station d’observation, puis les afficher exactement dans l’ordre des noms - Les implémentations doivent être faites uniquement en Java ; les distributions fournies par SDKMan et les builds Early Access d’openjdk.net sont autorisés, mais les dépendances externes sont interdites
- Les participants soumettent leur solution via une pull request sur le dépôt 1brc de GitHub, et peuvent comparer le format de réponse et les performances avec l’implémentation de base fournie
- L’évaluation se fait dans le même environnement Hetzner Cloud CCX33, avec 5 exécutions ; le classement du leaderboard est établi sur la moyenne de 3 exécutions après exclusion du meilleur et du pire temps
Le défi Java qui agrège le plus vite un milliard de lignes
- One Billion Row Challenge est un défi de performance Java qui s’est déroulé du 1er au 31 janvier 2024
- Les participants écrivent un programme Java qui lit des mesures de température dans un fichier texte et calcule les températures minimale, moyenne et maximale pour chaque station météo
- Le cœur de la difficulté tient au fait que le fichier d’entrée contient 1 000 000 000 de lignes
- L’entrée a une structure simple, avec une mesure par ligne
- Exemple :
Hamburg;12.0 - Exemple :
Bulawayo;8.9 - Exemple :
Palembang;38.8
- Exemple :
- La sortie doit trier les noms des stations par ordre alphabétique et afficher les valeurs
min/mean/maxde chaque station- Exemple :
{Abha=5.0/18.0/27.4, Abidjan=15.7/26.0/34.1, ...}
- Exemple :
Règles de soumission et environnement d’exécution
- L’objectif est de créer l’implémentation Java la plus rapide réalisant la même tâche
- Les optimisations peuvent exploiter les threads virtuels, la Vector API et SIMD, l’optimisation du GC, la compilation AOT, etc.
- Les règles de base sont les suivantes
- Les soumissions doivent être écrites en Java
- Les distributions Java proposées par SDKMan et les builds Early Access d’openjdk.net peuvent être utilisés
- Les builds EA de projets OpenJDK comme Valhalla sont également autorisés
- Les dépendances externes ne sont pas autorisées
- Les participants clonent le dépôt 1brc et soumettent leur implémentation en suivant les instructions du README
- L’implémentation de base est fournie comme point de comparaison et pour vérifier le format attendu de la réponse
- La soumission se fait en ouvrant une pull request vers le dépôt upstream
Méthode de calcul du leaderboard et partage communautaire
- L’évaluation est réalisée sur une instance Hetzner Cloud CCX33
- La configuration est de 8 vCPU dédiés et 32 Go de RAM
- Le temps d’exécution end-to-end est mesuré avec le programme
time - Chaque soumission est exécutée 5 fois de suite
- L’exécution la plus lente et la plus rapide sont exclues
- La moyenne des 3 temps d’exécution restants constitue le résultat de la soumission
- Le résultat est ajouté au leaderboard
- Les discussions sur les techniques d’optimisation se poursuivent dans les discussions du dépôt GitHub
- Un espace Show & Tell est également prévu pour partager des implémentations 1BRC dans d’autres langages que Java, comme Rust, Go ou C++
2 commentaires
Avis sur Hacker News
La solution [0] qui semble actuellement la plus performante ne tient pas compte des collisions de hachage ; si le jeu de données contient suffisamment de villes différentes, elle risque donc de produire des résultats incorrects.
Je me demande si quelque chose m’échappe.
[0] https://github.com/gunnarmorling/1brc/blob/main/src/main/jav...
Pour l’instant, ces entrées ont été retirées du classement, et les deux auteurs sont en train de corriger leurs soumissions ; elles devraient donc être réintégrées plus tard.
[0] https://twitter.com/mtopolnik/status/1742652716919251052
Avec l’approche suivante, je pense qu’on peut traiter l’ensemble en moins de 0,3 seconde.
Les températures n’ont qu’une décimale, donc environ 400 valeurs suffisent dans le cas général ; les noms de lieux sont eux aussi finis, autour de 400, ce qui permet de construire une table de correspondance d’environ 160 000 combinaisons température × lieu.
On génère automatiquement une machine à états qui mappe chacune de ces 160 000 combinaisons vers un bucket unique de la table de hachage, quelle que soit sa position de rotation dans un registre de 4 octets, puis, à chaque cycle, on effectue une recherche dans la table de transitions d’état et un XOR avec les 4 octets suivants dans un registre d’état 32 bits.
Il suffit de parcourir toutes les données à la vitesse de la mémoire et d’incrémenter les compteurs par état ; comme il n’y a que 65 K états, les compteurs tiennent dans le cache.
Avec AVX512, on peut faire tourner 512 de ces machines à états 32 bits en parallèle par cœur, donc le calcul ne devrait pas être le goulot d’étranglement.
Les températures trop hautes/basses ou les noms de lieux inconnus qui ne correspondent à aucun bucket valide sont envoyés vers du code lent, et le traitement des minimums/maximums peut aussi passer par cette échappatoire, ce qui ne se produit que quelques milliers de fois.
Cette méthode peut fonctionner à la vitesse de la mémoire avec un seul cœur AVX512, donc je ne pense pas qu’il y ait un intérêt à répartir sur plusieurs cœurs.
Il suffit d’une table de hachage de 400 entrées, de trois valeurs flottantes courantes pour minimum, moyenne et maximum, et d’un entier de comptage pour mettre à jour la moyenne.
Même avec 16 octets pour le nom, le tout tient dans moins de 16 Ko.
Le temps d’exécution sera dominé par les entrées/sorties, puis probablement par le parsing JSON.
La plupart des puces serveur x86 modernes peuvent retirer deux chargements SIMD par cycle ; avec AVX2, cela représente environ 32 Go/s à 1 GHz, donc AVX-512 n’est pas indispensable pour maximiser la bande passante par cœur.
Mais si la lecture se fait depuis la DRAM, on sera bloqué bien plus tôt, généralement autour de 10 à 16 Go/s sur un serveur.
Tant que la plupart des données débordent en RAM, le débit d’un seul cœur chute fortement, et pour les gros traitements en streaming, le parallélisme multicœur est presque toujours bénéfique.
C’est facile à vérifier : allouez un bloc mémoire bien plus gros que le cache L3, provoquez les défauts de page à l’avance, puis lancez dans une boucle serrée des chargements vectoriels déroulés (AVX2/AVX-512).
Je m’interroge aussi sur l’interprétation du registre d’état. Si on fait un XOR avec les 4 octets d’entrée, un nom de lieu inattendu peut en pratique donner n’importe laquelle des quelque 4,7 milliards de valeurs possibles.
Même pour un nom attendu, s’il fait plus de 4 octets, ne faut-il pas plusieurs états pour le distinguer d’autres noms partageant le même préfixe ?
Les règles disent que, même si le générateur de données utilise un ensemble fixe de noms de stations, toute solution doit fonctionner avec des noms de stations UTF-8 arbitraires.
Plutôt que d’écarter l’exécution la plus lente et la plus rapide puis de prendre la moyenne des trois restantes, je pense qu’il vaudrait mieux écarter les deux plus lentes, ou simplement retenir la valeur la plus rapide.
Je ne vois pas de bonne raison de jeter un bon résultat d’exécution.
Mais dès qu’il existe la moindre source de non-déterminisme à l’intérieur du programme — ce qui est plus fréquent qu’on ne le croit — le meilleur temps risque d’être peu représentatif.
À ce sujet, https://tratt.net/laurie/blog/2019/minimum_times_tend_to_mis... est un bon article.
Du point de vue de quelqu’un qui examine les règles à la lettre, on aurait envie de lancer un démon en arrière-plan au premier lancement, de charger tout le fichier en mémoire et de l’y épingler, puis de précharger le cache pour que les exécutions suivantes ne fassent pratiquement plus qu’un scan linéaire
Selon jusqu’où l’on étire l’interprétation des règles, précalculer le résultat dès le premier lancement semble aussi possible, et on pourrait même parser à l’avance les nombres dans un format plus dense puis les lire directement comme somme cumulée lors des exécutions suivantes
Ce n’est pas du tout dans l’esprit du concours, mais, d’après les règles visibles, cela ne semble pas interdit
Si l’on n’aime pas le précalcul, on peut aussi recourir à des astuces comme trier l’entrée à l’avance, la parser à l’avance, ou utiliser de la compression, du tri, ou une disposition mémoire triée
À l’extrême, on pourrait même patcher le script
calculate_timepour qu’il renvoie 0 seconde pour soi et 9999 pour les concurrentsEntre coder en dur la réponse sur une seule ligne sans même lire l’entrée, et traiter l’entrée en supposant qu’on n’en connaît pas le contenu, il existe environ un milliard de nuances dans la zone grise du précalcul
Cela peut devenir un concours visant à déterminer ce qui constitue un précalcul équitable ou non
C’est pourquoi les concours de machine learning ne montrent pas les données finales aux participants
Il est indiqué que le calcul doit avoir lieu au moment de l’exécution de l’application, et qu’il ne faut pas traiter le fichier mesuré au moment du build pour intégrer le résultat dans le binaire
J’ai l’impression que ce problème est simplement limité par la vitesse du disque. Je me demande si des optimisations comme SIMD ou le multithreading ont vraiment un intérêt
Cela dépendra certes du nombre de stations différentes et de la méthode de recherche dans la table de hachage, mais je doute que ce soit mesurable par rapport aux entrées-sorties
Les systèmes conçus pour du matériel moderne exploitent ce point, et redpanda.com, où je travaille, en est un exemple
Le parsing représente une grande part du temps de calcul, et des techniques SIMD comme SWAR pour trouver les séparateurs peuvent aider
Si vous voulez voir une implémentation propre de ce type d’algorithmes, Stringzilla est une bonne référence : https://github.com/ashvardanian/StringZilla
J’ai répondu ici au fait que le fichier soit entièrement mis en cache en mémoire après la première exécution : https://news.ycombinator.com/item?id=38864034
Un serveur typique dispose d’assez de lignes PCIe pour accueillir 15 SSD de ce type, donc la bande passante d’E/S d’un serveur est du même ordre que la bande passante mémoire
Les serveurs plus chers ont davantage de lignes, et plus rapides, comme du PCIe 5.0
Ce fichier contient un milliard de lignes et fait environ 1 Go compressé ; après la première exécution écartée, il tient en mémoire, donc dans ce scénario la bande passante d’E/S n’est pas importante
Le dépôt GitHub indique 12 Go non compressés, ce qui confirme malgré tout que la bande passante d’E/S n’est pas déterminante
L’idée principale est que le disque est rarement le goulot d’étranglement
Par exemple, sous Linux avec ext2, il y a de fortes chances que tout le fichier soit mis en cache après la première exécution, mais ce n’est pas forcément le cas avec ZFS
Ainsi, les chiffres arrivent des unités vers les positions les plus significatives, puis viennent le séparateur et la chaîne, jusqu’à rencontrer EOF ou un saut de ligne
D’après les règles, une soumission doit fonctionner correctement pour toutes les entrées, mais elle peut — et devrait probablement — être optimisée pour l’entrée spécifique générée par
create_measurements.shPar exemple, on peut imaginer une soumission utilisant une fonction de hachage parfaite adaptée à l’ensemble de stations fourni
Cela permettrait d’éviter les optimisations par surapprentissage
Un octet supérieur à 127 signifie un caractère UTF-8 multioctet
Par curiosité, j’ai comparé les performances de awk contre Java
Il s’agit d’un script qui, avec
awk -F';', cumule par station la somme, le nombre, le minimum et le maximum, puis calcule et affiche la moyenne dans ENDL’idée serait de créer une table externe sur le fichier CSV avec
file_fdw, puis de calculerMIN,AVGetMAXavecGROUP BY station_nameDans
clickhouse local, on litfile('measurements.txt', 'CSV', 'station String, t Float32'), on groupe par station avecmin,maxetavg, et on exécute avecmax_threads = 8La majeure partie du temps est passée à parser le fichier
sumpeut devenir assez grande, il vaut mieux utiliser une moyenne en streamingPar exemple, une formule du type
new_mean = ((n*old_mean)+temp)/(n+1)Un défi intéressant, mais dommage qu’il soit réservé à Java. J’attends avec impatience le moment où les gens commenceront à écrire eux-mêmes du bytecode JVM à la main
[0] https://github.com/gunnarmorling/1brc/discussions
C’est amusant. Ça fait un peu after d’Advent of Code
Pour une comparaison équitable entre langages, il faudrait aussi inclure
makeet le temps de build. Je n’ai pas utilisé Java/Maven depuis quelques années, mais en voyant./mvnw clean verifytélécharger des choses depuis 2 minutes, je me rappelle pourquoiEt pour les outils de build à compilation incrémentale, Gradle est plus rapide
Il faudrait aussi ajouter une part appropriée du temps nécessaire pour apprendre à programmer
Dans ce genre de défi, une version très naïve aurait de bonnes chances de gagner ; à mon avis, ce serait non seulement irréaliste, mais aussi contraire à l’esprit du défi
cleanC’est comme jeter le cache, puis dire que c’est lent
Il est indiqué qu’on ne peut pas utiliser de dépendances externes
À l’Université technique de Prague, dans un cours de C, il y avait un exercice très similaire
Toutes les soumissions des étudiants étaient évaluées en continu dans un classement, et beaucoup d’étudiants passaient des dizaines d’heures à optimiser pour obtenir des points bonus — en pratique, des points de prestige — afin d’avoir une meilleure note
Le premier est à 6 secondes… impressionnant.