Skip to Content
Mélodium 0.10.3 is now available!
DocsExemplesCluster LLM distribué

Cluster LLM distribué

Source: showcase/distributed_llm_cluster See in Playground

Il s’appuie directement sur les primitives distrib de l’étape tutoriel sur le calcul distribué, en ajoutant work/distant pour provisionner un moteur distant à la demande depuis Cadence.CI, plutôt que de pointer vers un nœud melodium dist démarré manuellement. Un serveur HTTP accepte des requêtes POST /chat avec un prompt en texte brut et retransmet les tokens générés au fur et à mesure. Cadence.CI provisionne un worker qui télécharge un modèle Mistral depuis le HuggingFace Hub et le charge en mémoire, une seule fois ; le paquet ml, les poids téléchargés et la puissance de calcul nécessaire pour les exécuter n’ont besoin d’être disponibles que sur le worker lui-même, jamais sur le processus côté serveur.

Note

Nécessite un jeton d’API Cadence.CI. Placez MELODIUM_API_TOKEN dans l’environnement et exécutez avec --api-report pour suivre l’exécution là-bas. Aucune clé d’API d’un fournisseur de LLM n’est nécessaire : l’inférence s’exécute sur le worker provisionné lui-même, sur un modèle Mistral téléchargé depuis le HuggingFace Hub, et non contre une API tierce hébergée.

Exécution

cd showcase/distributed_llm_cluster export MELODIUM_API_TOKEN="my-cadence-ci-token" melodium run --api-report Compo.toml --port 8080 curl -X POST http://127.0.0.1:8080/chat \ -d "Explain the Mélodium dataflow model in one sentence."

Fonctionnement

Un modèle DistantEngine demande un worker à Cadence.CI, un modèle DistributionEngine nomme le traitement distant (inferText) qui s’y exécute, et un modèle HttpServer sert d’écouteur côté serveur :

model runner: DistantEngine(api_token=_, api_url=_) model distributor: DistributionEngine( treatment = "distributed_llm_cluster/main::inferText", version = "0.1.0" ) model httpServer: HttpServer(host=|from_ipv4(|localhost_ipv4()), port=port)

distant demande un worker dimensionné pour une véritable inférence locale (16 Go de mémoire, 4 CPU, 32 Go de stockage, suffisant pour les poids fp16 de Mistral-7B plus la marge pour le moteur et le cache du Hub), puis distrib::start s’y connecte, en transmettant l’identifiant du dépôt HuggingFace une seule fois comme entrée params plutôt qu’à chaque requête :

provisionRunner: distant[distant_engine=runner]( max_duration = 3600, memory = 16384, cpu = 4000, storage = 32768, edition = _, arch = _, volumes = [], containers = [], service_containers = [], tags = [] ) startup.trigger -> provisionRunner.trigger,access -> connectDistributor.access connectDistributor: distribStart[distributor=distributor](params=|dataMap([|dataEntry<string>("repo_id", hf_repo_id)]))

Une connexion distrib établie signifie seulement que le processus worker est joignable, pas que ses poids de modèle, plusieurs gigaoctets, ont fini d’être téléchargés et chargés, ce qui peut prendre plusieurs minutes la première fois. Le processus côté serveur envoie une véritable requête de préchauffage juste après la connexion, et ne démarre le serveur HTTP qu’une fois cette requête effectivement répondue :

warmupPrompt: emit<string>(value=" ") connectDistributor.ready -> warmupPrompt.trigger warmupStream: stream<string>() warmupPrompt.emit -> warmupStream.block warmupBytes: encode() warmupStream.stream -> warmupBytes.text warmupCall: dispatchInfer[distributor=distributor]() warmupBytes.data -> warmupCall.prompt warmupDone: trigger<byte>() warmupCall.response -> warmupDone.stream startHttp: start[http_server=httpServer]() warmupDone.start -> startHttp.trigger

Côté worker, release retient chaque prompt entrant, y compris celui de préchauffage, derrière le signal loaded du modèle avant qu’il n’atteigne generate : un prompt qui arrive pendant le chargement attend simplement là plutôt que d’être écarté. C’est ce qui transforme la réponse de préchauffage en un signal de disponibilité fiable, plutôt qu’une course contre un modèle encore en cours de chargement.

Distribuer une requête vers le worker

dispatchInfer suit la même poignée de main en trois étapes que l’étape tutoriel sur le calcul distribué : allouer un identifiant de distribution, envoyer l’entrée, recevoir la sortie, le tout étiqueté par nom. Il porte à la fois la requête de préchauffage ci-dessus et chaque véritable requête de chat :

treatment dispatchInfer[distributor: DistributionEngine]() input prompt: Stream<byte> output response: Stream<byte> { trig: trigger<byte>() dist: distribute[distributor=distributor]() Self.prompt -> trig.stream,start -> dist.trigger sendPrompt: sendStream<byte>[distributor=distributor](name="prompt") recvResponse: recvStream<byte>[distributor=distributor](name="response") dist.distribution_id -> sendPrompt.distribution_id dist.distribution_id -> recvResponse.distribution_id Self.prompt -> sendPrompt.data recvResponse.data -> Self.response }

Récupérer et charger le modèle une seule fois

inferText s’exécute sur le worker, lancé une seule fois quand distribStart s’y connecte plutôt qu’à chaque requête : son propre startup() se déclenche à ce moment-là. Il récupère les poids du modèle depuis le HuggingFace Hub avec un modèle HfHub et fetch, puis les charge dans un modèle Mistral avec load une seule fois, et sert ensuite chaque requête de chat contre ce même modèle chargé avec generate :

treatment inferText(const repo_id: string) model hub: HfHub(repo_id=repo_id) model mistral: Mistral(max_new_tokens=256) input prompt: Stream<byte> output response: Stream<byte> { startup() fetchWeights: fetch[hub=hub]() loadModel: load[mistral=mistral]() startup.trigger -> fetchWeights.trigger fetchWeights.safetensors -> loadModel.safetensors fetchWeights.tokenizer -> loadModel.tokenizer waitForModel: release<byte>() loadModel.loaded -> waitForModel.leverage Self.prompt -> waitForModel.data decodePrompt: decode() waitForModel.released -> decodePrompt.data,text -> generateReply.prompt generateReply: generate[mistral=mistral]() encodeResponse: encode() generateReply.generated -> encodeResponse.text,data -> Self.response }

Dépendances

[dependencies] std = "0.10.3" # flux de base, journalisation, structures de données http = "0.10.3" # client et serveur HTTP net = "0.10.3" # utilitaires d'adresses IP encoding = "0.10.3" # encodage / décodage UTF-8 work = "0.10.3" # provisionnement de runners cloud distrib = "0.10.3" # distribution de flux entre runners ml = "0.10.3" # LLM, STT, TTS et inférence de modèles locaux