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