Paylanmış MCMC Konveyerlərinin Arxitekturası: Go gRPC, pgvector və Next.js Server-Sent Events
Go-da paralel nümunələmə, gRPC üzərindən binar daşıma, pgvector-də posterior vəziyyət və brauzerə axınla ötürülən konvergensiya metrikləri.
Yerli mühitdə R və ya Python ilə işləyən yüksək optimallaşdırılmış statistik skriptdən istehsalat (production) səviyyəli, paylanmış veb arxitekturasına keçid tətbiqi məlumat elmində (data science) ən kritik mühəndislik maneələrindən biridir. Markov Zənciri Monte Karlo (MCMC) alqoritmləri, məsələn, Hamiltonian Monte Karlo (HMC) və ya adaptiv Metropolis-Hastings, təbiəti etibarilə ardıcıl, hesablama baxımından çox ağır və davamlı yaddaş ayrılmasına (memory allocation) ciddi şəkildə ehtiyac duyan proseslərdir. Məlumat elmi komandaları bu modelləri standart Python (məsələn, FastAPI və ya Flask) və ya R (məsələn, Plumber) mühitlərindən istifadə edərək monolit REST API-ləri kimi yerləşdirməyə cəhd etdikdə, arxitektura qaçılmaz olaraq paralel sorğuların (concurrent load) təzyiqi altında çökür. Bu icra mühitlərinin sinxron təbiəti icra axınlarını (threads) bloklayır, nəticədə böyük həcmdə yaddaş sərfiyyatına, kəsilən bağlantılara və gecikmənin (latency) fəlakətli dərəcədə pisləşməsinə səbəb olur.
2026-cı il etibarilə, ehtimal modellərinin yerləşdirilməsi üçün sənaye standartı tamamilə bir-birindən asılı olmayan (decoupled) və başlıqsız (headless) arxitekturanı tələb edir. Biz ağır riyazi nüvələri şəbəkə daşınması təbəqəsindən ayırmalıyıq. Müasir paradiqma statistik alqoritmləri yüksək asinxronluğa malik Go mikroxidmətlərində kapsullaşdırmağı, yaradılan posterior nümunələrini optimallaşdırılmış gRPC kanalları vasitəsilə yayımlamağı, vəziyyət vektorlarını pgvector istifadə edərək PostgreSQL-də yaddaşa yazmağı və konvergensiya metriklərini Next.js App Router ön üzündə (frontend) Server-Sent Events (SSE) vasitəsilə qəbul etməyi nəzərdə tutur. Bu hesabat həmin sistemin qurulmasına aparan tam mühəndislik yolunu təfərrüatı ilə izah edir.
Paylanmış MCMC Mühərriki: Go-da Asinxronluq
Mürəkkəb bir posterior paylanmanı dəqiq qiymətləndirmək üçün geniş səpələnmiş başlanğıc nöqtələrindən başlayan sayda paralel Markov zəncirini icra etmək standart təcrübədir. Ənənəvi böyük məlumat (big data) mühitlərində, bu zəncirləri məlumat seqmentləri (shards) üzrə paylamaq üçün adətən PySpark kimi çərçivələrdən istifadə olunur. Bu yanaşma işləyir və onu rədd etmək əvəzinə nə verdiyini dəqiq demək daha faydalıdır: PySpark üzərində paylanmış Hamiltonian Monte Karlo üzrə 2025-ci il tədqiqatı sətir və ölçülü sintetik loqistik reqressiyada, 4-dən 32-yə qədər işçi ilə, 0.986 qəbul dərəcəsində saniyədə təxminən 18.7 Effektiv Nümunə Ölçüsü (ESS/s) bildirir. Həmin tədqiqatda paylanmış baza variantı müqayisə edildiyi kommunikasiyadan qaçan variantı geridə qoymuşdu.
Bu rəqəm faydalı istinad nöqtəsidir, universal tavan deyil, və onu əvəzləmək iddiasında olan hər arxitektura məhz bu göstərici ilə ölçülməlidir. Daşıdığı xərclər isə strukturaldır: hər icraçıya düşən JVM yükü, paylanmış fayl sisteminin gecikməsi və hər model yeniləməsindəki qlobal sinxronizasiya baryeri. Bunların heç biri nümunə toplamanın öz təbiətindən doğmur, məhz buna görə də mühəndislik yolu ilə aradan qaldırılmağa dəyər.
Nümunə toplama orkestrasiyasının Go-da yenidən yazılması bu məhdudiyyətləri aradan qaldırır. Go-nun yüngül "goroutine"-ləri və kilidsiz (lock-free) kanalları, əməliyyat sistemi səviyyəli axınlara (OS-level threads) xas olan ağır kontekst dəyişməsi (context-switching) cəzası olmadan tək bir çoxnüvəli maşında minlərlə müstəqil Markov zəncirinin icrasına imkan verir.
gRPC Axını və vtprotobuf Optimallaşdırmaları
Statistik mühərrik ilə veb tətbiqetmə təbəqəsi arasında şəbəkə daşınması üçün HTTP/1.1 üzərindən REST struktur olaraq qeyri-kafidir. REST JSON serializasiyasına əsaslanır ki, bu da hesablama baxımından bahalı olan sətir (string) analizini tələb edir və yaddaş təmizləyicisi (GC) üzərində əhəmiyyətli təzyiq yaradır. Əvəzində, HTTP/2 üzərindən gRPC, binar serializasiya üçün Protocol Buffers (protobuf) istifadə edərək multipleksləşdirilmiş, davamlı bağlantılar təmin edir.
gRPC istifadə etmək qərarı birbaşa aparat təminatının (hardware) istifadə metriklərindən irəli gəlir. Standart protobuf serializasiyası sürətlidir, lakin MCMC çıxışının miqyasında, bir xidmətin saniyədə on minlərlə çoxölçülü sürüşən nöqtəli (float) massiv yayımlaya bildiyi halda, yaddaşın ayrılması (memory allocation) əsas maneəyə çevrilir. vtprotobuf qoşmasının (plugin) tətbiq edilməsi daxil olan və xaric olan axınlar (ingress/egress streams) zamanı yaddaş ayırmalarını hovuzlaşdıraraq (memory pooling) və təkrarlanan struktur (struct) yaradılmasının qarşısını alaraq bu prosesi kəskin şəkildə optimallaşdırır.
| Metrik | REST (JSON / HTTP/1.1) | Standart gRPC (Protobuf) | Optimallaşdırılmış gRPC (vtprotobuf + Yaddaş Hovuzu) |
|---|---|---|---|
| Serializasiya Formatı | Mətn (JSON) | Binar | Binar |
| Bağlantı Yükü | Yüksək (Hər Sorğu Üçün) | Aşağı (Multipleksləşdirilmiş) | Aşağı (Multipleksləşdirilmiş) |
| CPU İstifadəsi (Xaric olan Axın) | ~3.0 Nüvə | ~1.6 Nüvə | ~0.7 Nüvə |
| Yaddaş Ayrılması | Davamlı GC təzyiqi | Orta | Minimal (Hovuzlaşdırılmış) |
| MCMC üçün Uyğunluğu | Zəif | Yaxşı | Mükəmməl |
Aşağıda "zombi" hesablamaların yaddaş sızıntısına səbəb olmasının qarşısını almaq üçün müştəri ayrılmalarını təhlükəsiz şəkildə idarə edən və paralel zəncirləri orkestr etmək üçün kanallardan istifadə edən istehsalat səviyyəli Go gRPC xidmətinin tətbiqi verilmişdir.
package main
import (
"context"
"math/rand"
"sync"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
pb "github.com/your-org/mcmc-core/proto" // vtprotobuf generasiyasını fərz edir
)
type MCMCServer struct {
pb.UnimplementedMCMCServiceServer
}
// StreamPosterior paralel zəncirləri icra edir və nümunələri koordinatora axınla ötürür.
func (s *MCMCServer) StreamPosterior(
req *pb.SamplingRequest,
stream pb.MCMCService_StreamPosteriorServer,
) error {
var wg sync.WaitGroup
// Şəbəkədəki qısamüddətli əks-təzyiqi udmaq üçün buferləşdirilmiş kanal.
sampleChan := make(chan *pb.Sample, 5000)
// Bu kontekstin ləğvi hər zəncirə dayanmağı bildirən siqnaldır. O, həm
// müştəri bağlantını kəsdikdə stream.Context() vasitəsilə, həm də biz özümüz
// geri döndükdə aşağıdakı defer cancel ilə işə düşür.
ctx, cancel := context.WithCancel(stream.Context())
defer cancel()
// Paralel Markov zəncirlərini başlat.
for i := int32(0); i < req.NumChains; i++ {
wg.Add(1)
go func(chainID int32) {
defer wg.Done()
// Dövr daxilində yaddaş ayırmalarını minimuma endirmək üçün vəziyyət
// massivini əvvəlcədən ayır.
state := make([]float64, req.Dimensions)
for step := int32(0); step < req.Iterations; step++ {
// Metropolis-Hastings təklif addımını simulyasiya et.
for d := int32(0); d < req.Dimensions; d++ {
state[d] += rand.NormFloat64() * 0.1
}
sample := &pb.Sample{
ChainId: chainID,
Iteration: step,
State: append([]float64(nil), state...),
}
// Göndərmə select-in içində olmalıdır, ondan sonra yox. Ayrıca
// `select { case <-ctx.Done(): ... default: }` ilə qorunan sadə
// `sampleChan <- sample` yazılışında bağlantını kəsmiş müştəri
// bütün istehsalçıları dolu buferin qarşısında əbədi saxlayır:
// wg.Wait heç vaxt geri dönmür, kanal bağlanmır və goroutine-lər
// sızır. Bu blokun mövcud olma səbəbi olan "zombi" hesablama məhz
// budur, ona görə ləğvetmə göndərmənin özünü də əhatə etməlidir.
select {
case sampleChan <- sample:
case <-ctx.Done():
return // Müştəri ayrıldı; zənciri tərk et.
}
}
}(i)
}
// Bütün zəncirlər bitdikdən sonra kanalı bağla.
go func() {
wg.Wait()
close(sampleChan)
}()
// Axın istehlakı dövrü. Buradan geri dönmək defer edilmiş cancel-i işə salır,
// o da yuxarıdakı göndərmədə hələ gözləyən istehsalçıları azad edir.
for sample := range sampleChan {
if err := stream.Send(sample); err != nil {
return status.Errorf(codes.Unavailable, "axına göndərmə uğursuz oldu: %v", err)
}
}
return nil
}
Çoxdilli MCMC Konveyerinin Dockerləşdirilməsi
Bu arxitekturanı yerləşdirmək üçün Go binar faylı konteynerləşdirilməlidir. Köhnə R və ya C++ statistik kitabxanalarını (məsələn, RStan və ya xüsusi Hamiltonian tətbiqləri) cgo vasitəsilə inteqrasiya edərkən, Docker qurulma prosesi son istehsalat obrazının (image) yüngül və təhlükəsiz qalmasını təmin etmək üçün çoxmərhələli (multi-stage) bir yanaşma tələb edir.
İlk mərhələ lazımi C++ alətlərini və R başlıqlarını ehtiva edən ağır kompayler obrazından istifadə edir. İkinci mərhələdəki həlledici məhdudiyyət isə C kitabxanasıdır. Debian əsaslı qurucuda CGO_ENABLED=1 ilə kompilyasiya edilmiş binar fayl glibc-yə bağlanır, Alpine isə musl daşıyır, ona görə həmin faylı Alpine obrazına köçürmək işə düşən kimi "yükləyici tapılmadı" xətası verən konteyner yaradır. Bundan başqa, R-a olan cgo bağlantıları statik əlaqələnmir: libR.so paylaşılan obyektdir və icra obrazı onu özü ilə daşımalıdır. Təmiz Go xidmətləri üçün gəzən alpine üstəgəl -extldflags "-static" resepti R işin içinə girən kimi sadəcə tətbiq olunmur.
Düzgün cütlük qurucu ilə eyni libc üzərində, R icra mühitini daşıyan və başqa heç nə saxlamayan obrazdır:
# Mərhələ 1: qurulma mühiti
FROM golang:1.24-bookworm AS builder
# cgo bağlantılarının kompilyasiya olunduğu C++ alətləri və R başlıqları.
RUN apt-get update && apt-get install -y --no-install-recommends \
build-essential \
r-base-dev \
r-cran-rcpp \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
# Qəsdən dinamik əlaqələndirilir. libR Go binar faylının içinə statik
# yerləşdirilə bilməz, ona görə aşağıdakı icra obrazı onu təmin etməlidir.
RUN CGO_ENABLED=1 GOOS=linux go build -o mcmc_engine .
# Mərhələ 2: eyni libc üzərində icra obrazı, R var, alətlər yoxdur.
FROM debian:bookworm-slim
RUN apt-get update && apt-get install -y --no-install-recommends \
ca-certificates \
r-base-core \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /app
COPY --from=builder /app/mcmc_engine .
EXPOSE 50051
CMD ["./mcmc_engine"]
Bu, obrazı çoxgiqabaytlıq qurucudan təxminən 400 MB-a endirir: kompilyator, başlıq faylları və R inkişaf paketləri gedir, R icra mühiti isə qalır, çünki binar fayl ona həqiqətən möhtacdır. Go xidmətləri üçün deyilən 50 MB-dan aşağı rəqəmə yalnız cgo-nu tənlikdən çıxarmaqla çatmaq olar və sampler-i təmiz Go-da yazmaq mümkün olan hər yerdə bunu etməyə dəyər. CGO_ENABLED=0 ilə binar tam statik olur, icra mərhələsi FROM scratch üstəgəl CA sertifikatlarına çevrilir və obraz təxminən 20 MB-a düşür. R bağlantılarına daimi arxitektura kimi yox, keçid mərhələsi kimi baxın.
Dinamik Konvergensiya Diaqnostikası ()
Avtomatlaşdırılmış MCMC yerləşdirilməsində ən böyük risk prosesin vaxtından əvvəl dayandırılmasıdır. Zəncirləri hədəf stasionar paylanmaya çatmazdan əvvəl dayandırmaq güclü qərəzli (biased) parametr qiymətləndirmələri ilə nəticələnir, zəncirləri sonsuzadək işlətmək isə hesablama resurslarını israf edir və infrastruktur xərclərini artırır. Stasionarlığı müəyyən etmək üçün sənaye standartı mexanizmi Gelman-Rubin potensial miqyasın kiçilməsi faktorudur, yəni .
-in riyazi əsası çoxsaylı müstəqil zəncirlər arasındakı dispersiyanı () həmin fərdi zəncirlərin daxilindəki dispersiya () ilə müqayisə etməyə əsaslanır. Əgər zəncirlər eyni posterior fəzasına konvergensiya edibsə (yaxınlaşıbsa), zəncirlərarası dispersiya və zəncirdaxili dispersiya demək olar ki, eyni olmalıdır.
Fərz edək ki, zəncirlərin sayı, isə hər zəncir üçün iterasiyaların sayıdır. Zəncirlərarası dispersiya belə müəyyən edilir:
Zəncirdaxili dispersiya belə müəyyən edilir:
Birləşdirilmiş marjinal dispersiya qiymətləndirməsi çəkili orta dəyər kimi hesablanır:
Sonra potensial miqyasın kiçilməsi faktoru aşağıdakı kimi alınır:
Yuxarıdakı formula 1992-ci ilin klassik diaqnostikasıdır və ilk tətbiq ediləcək variant məhz odur, çünki tənliklərin təsvir etdiyi versiya budur. Onun avtomatlaşdırılmış dayandırma qaydasında əhəmiyyət kəsb edən iki kor nöqtəsi var. Ortalaması sabit, lakin dispersiyası sürüşən zənciri görmür, həmçinin hədəfin sonlu ortalama və dispersiyaya malik olduğunu fərz edir ki, bu da ağır quyruqlu posteriorlarda pozulur. Hər ikisi Vehtari və həmmüəlliflərinin (2021) rank-normallaşdırılmış, qatlanmış split- variantı ilə həll olunur: hər zənciri iki yerə bölmək zəncirdaxili trendi zəncirlərarası fərqə çevirir, rank normallaşdırma isə moment fərziyyəsini aradan qaldırır. R-dakı posterior::rhat() və Python-dakı arviz.rhat bu gün məhz bunu hesablayır və istehsalat orkestratoru öz əli ilə yazılmış formul əvəzinə onlardan birini çağırmalıdır.
Hədd də diaqnostika ilə birlikdə dəyişdi. Uzun müddət qəbul edilən o qədər yumşaq idi ki, marjinal dispersiyanın 30%-nə qədərini təşkil edən trendləri görmürdü; Vehtari və həmmüəllifləri 2% səviyyəsindəki trendləri tutan həddini tövsiyə edir. Aşağıdakı dayandırma qaydası məhz buna görə 1.01 istifadə edir. Daha yumşaq hər bir dəyər hesablama müqabilində qərəzi qəbul etmək qərarıdır; bu, qanuni bir güzəştdir, sadəcə on illik konvensiyadan miras qalmaq yox, şüurlu şəkildə verilməlidir.
İstehsalat mühitində bu diaqnostika geriyə dönük hesablana bilməz. O, daxil olan gRPC axını üzərində dinamik olaraq hesablanmalıdır. Axın irəlilədikcə həm , həm də üçün kvadratların davamlı cəmi yenilənir və meyar bütün ölçülər üzrə ödəndikdə orkestrator zəncirlərin qarışdığını iddia edə bilər, gRPC axınını sonlandırmaq və goroutine-ləri azad etmək üçün erkən çıxış siqnalını işə salır.
PostgreSQL 16 və pgvector ilə Çoxölçülü Vəziyyətin Saxlanması
Davamlı posterior vəziyyət axınını saxlamaq çətindir. Tipik bir icra prosesi 4 zəncir üzrə 10,000 iterasiya yarada bilər və hər bir vəziyyət yüzlərlə ölçüdən ibarət olur. Tarixən bu məlumatlar düzəldilərək HDFS-ə və ya obyekt yaddaşına (blob storage) yazılırdı. Lakin pgvector genişlənməsi (v0.8+) ilə birləşdirilmiş PostgreSQL 16, ehtimal paylanmalarını (probabilistic distributions) necə saxladığımızı və təhlil etdiyimizi kökündən dəyişdirir.

MCMC vəziyyət vektorlarını VECTOR tiplərindən istifadə edərək yerli olaraq saxlamaqla biz mürəkkəb analitikaları birbaşa verilənlər bazasının daxilində yerinə yetirə bilərik. Məsələn, posterior nümunələrə k-ən yaxın qonşu (k-NN) sorğularını tətbiq etmək, məlumat alimlərinə çoxmodlu paylanmaları sürətlə aşkar etməyə və ya zaman keçdikcə zəncirlərin trayektoriyasının necə klasterləşdiyini təhlil etməyə imkan verir.
Milyonlarla nümunə üzrə yüksək sürətli sorğuları dəstəkləmək üçün vektorlara Hierarchical Navigable Small World (HNSW) indeksi tətbiq edilir.
| İndeks Parametri | Dəyər | MCMC Posterior Məlumatları Üçün Nəticələri |
|---|---|---|
| İndeks Tipi | hnsw |
Dinamik axınlar üçün IVFFlat-dan üstündür; əvvəlcədən təlim mərhələsi tələb etmir. |
| Məsafə Metrikası | vector_l2_ops |
Evklid məsafəsi posterior daxilindəki parametr vəziyyətləri arasındakı fəza məsafəsini dəqiq təmsil edir. |
| m | 16 |
İndeksin qurulması zamanı hər bir element üçün yaradılan ikitərəfli əlaqələrin maksimum sayını idarə edir. 16 yadda tutma dərəcəsini (recall) daxiletmə sürəti ilə tarazlaşdırır. |
| ef_construction | 64 |
Qrafın qurulması zamanı dinamik namizəd siyahısının ölçüsünü müəyyən edir. Daha yüksək dəyərlər daha yavaş daxiletmə bahasına sorğunun dəqiqliyini yaxşılaşdırır. |
Tələb olunan verilənlər bazası sxemi yüksək səviyyədə normallaşdırılmış SQL strukturu vasitəsilə icra edilir ki, bu da işin (job) metadatası üçün tranzaksiya bütövlüyünü təmin etməklə yanaşı nümunələr üçün vektor axtarış imkanları verir:
-- Çoxölçülü dəstək üçün pgvector genişlənməsini aktivləşdir.
CREATE EXTENSION IF NOT EXISTS vector;
-- MCMC icraları üçün ciddi relasiyalı bütövlüyü təmin edən metadata cədvəli.
CREATE TABLE mcmc_jobs (
job_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
model_name VARCHAR(255) NOT NULL,
dimensions INT NOT NULL,
status VARCHAR(50) DEFAULT 'RUNNING',
created_at TIMESTAMPTZ DEFAULT NOW()
);
-- Fərdi posterior vektorları saxlayan zaman seriyası cədvəli.
CREATE TABLE posterior_samples (
sample_id BIGSERIAL PRIMARY KEY,
job_id UUID REFERENCES mcmc_jobs(job_id) ON DELETE CASCADE,
chain_id INT NOT NULL,
iteration INT NOT NULL,
-- `vector` tipi 16,000 ölçüyə qədər saxlayır, lakin HNSW və ya IVFFlat
-- indeksi bu tip üçün 2,000 ilə məhduddur (`halfvec` üçün 4,000).
-- Parametr fəzasını saxlama limitinə görə deyil, indeks limitinə görə seçin.
state_vector VECTOR(256),
log_likelihood FLOAT,
created_at TIMESTAMPTZ DEFAULT NOW()
);
-- L2 məsafə məntiqindən istifadə edərək HNSW indeksini qur.
CREATE INDEX idx_posterior_hnsw
ON posterior_samples
USING hnsw (state_vector vector_l2_ops)
WITH (m = 16, ef_construction = 64);
RSC Sərhədini Aşmaq: Next.js Server-Sent Events
Sonuncu arxitektura komponenti real vaxt konvergensiya məlumatlarının vizuallaşdırma üçün müştərinin brauzerinə çatdırılmasını əhatə edir. REST son nöqtəsini dövri olaraq yoxlamaq (polling) olduqca səmərəsizdir və WebSockets tam dupleks əlaqə təklif etsə də, bağlantı vəziyyətinin idarə edilməsi, təhlükəsizlik divarından (firewall) keçid və yük tarazlayıcı (load balancer) konfiqurasiyası baxımından ciddi infrastruktur mürəkkəbliyi yaradır.
Əvəzində, Server-Sent Events (SSE) mükəmməl birtərəfli ötürmə mühiti təmin edir. Next.js App Router-dən istifadə edərək, yerli veb ReadableStream API vasitəsilə SSE son nöqtəsi (endpoint) qururuq. Bu texnika Next.js-ə standart HTTP keep-alive bağlantısı üzərindən məlumat parçalarını (chunks) serverdən müştəriyə tədricən ötürməyə imkan verir.
Route Handler bunun üçün məhz React render yolundan kənarda dayandığına görə düzgün yerdir. Server Komponentləri bir dəfə render olunub bitir; MCMC icrası boyunca açıq qalan axının həmin həyat dövründə yeri yoxdur. Bağlantını Route Handler-də saxlamaq və onu müştəri komponentində EventSource ilə oxumaq deməkdir ki, axının vəziyyəti tamamilə müştəridə yaşayır və server renderində heç vaxt iştirak etmir, yəni server ilə müştəri ağaclarının razılaşmayacağı bir şey qalmır. Bütün yolun izlənməsi isə ayrı məsələdir: Go gRPC çağırışından Next.js idarəedicisinə qədər ötürülən OpenTelemetry izi gecikmənin şəbəkədən yox, sampler-dən qaynaqlandığını müəyyən etməyə imkan verir.
// app/api/mcmc-stream/route.ts
import { NextRequest } from "next/server";
import { createGrpcClient, createDbPool } from "@/lib/infrastructure";
import { calculateRHat } from "@/lib/statistics";
import type { Sample } from "@/lib/proto/mcmc";
// Dinamik icranı məcbur et; axının statik keşlənməsinin qarşısını al.
export const dynamic = "force-dynamic";
export async function GET(req: NextRequest) {
const encoder = new TextEncoder();
const jobId = req.nextUrl.searchParams.get("jobId");
const stream = new ReadableStream({
async start(controller) {
const grpcClient = createGrpcClient();
const db = createDbPool();
const stateBuffer: Sample[] = [];
let pending: Sample[] = [];
// Hər nümunə üçün ayrıca INSERT bütün axını verilənlər bazasına gediş-gəliş
// müddətinə bağlayardı və Go mühərrikinin qurulma məqsədi olan ötürücülüyü
// yerə atardı. Əvəzində nümunələr yığılır və tək çoxsətirli ifadə ilə
// boşaldılır.
const FLUSH_EVERY = 500;
async function flush() {
if (pending.length === 0) return;
const values: unknown[] = [];
const tuples = pending.map((s, i) => {
const o = i * 4;
values.push(jobId, s.chainId, s.iteration, `[${s.state.join(",")}]`);
return `($${o + 1}, $${o + 2}, $${o + 3}, $${o + 4}::vector)`;
});
await db.query(
`INSERT INTO posterior_samples (job_id, chain_id, iteration, state_vector)
VALUES ${tuples.join(", ")}`,
values
);
pending = [];
}
try {
const streamCall = grpcClient.streamPosterior({
numChains: 4,
dimensions: 256,
iterations: 10000,
});
for await (const sample of streamCall) {
pending.push(sample);
stateBuffer.push(sample);
if (pending.length >= FLUSH_EVERY) await flush();
// Dövri olaraq Gelman-Rubin konvergensiya diaqnostikasını hesabla.
if (stateBuffer.length % 100 === 0) {
const rHat = calculateRHat(stateBuffer);
const data = JSON.stringify({
chain: sample.chainId,
iteration: sample.iteration,
primary_dimension: sample.state[0], // İzləmə qrafiki üçün çıxarılan xüsusiyyət.
rHat,
});
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
// Dayandırma qaydası: Vehtari və həmmüəlliflərinin (2021)
// rank-normallaşdırılmış split-R-hat həddi, köhnə 1.1 konvensiyası yox.
if (rHat < 1.01 && sample.iteration > 2000) {
controller.enqueue(
encoder.encode(`event: convergence\ndata: {"status": "converged"}\n\n`)
);
await flush();
await db.query(
`UPDATE mcmc_jobs SET status = 'COMPLETED' WHERE job_id = $1`,
[jobId]
);
break; // gRPC axınını zərifliklə bağla.
}
}
}
await flush(); // Son natamam paketdən qalan nümunələr.
} catch (error) {
console.error("SSE axını xətası:", error);
controller.enqueue(
encoder.encode(`event: error\ndata: {"message": "Nəticəçıxarma uğursuz oldu"}\n\n`)
);
} finally {
controller.close();
}
},
});
// Standart SSE başlıqlarını brauzerə qaytar.
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
// Bu olmadan Nginx və əksər edge proksiləri cavabı tamponlayır və bütün
// hadisələri idarəedici bitəndə birdəfəyə çatdırır; uzunömürlü axın üçün
// bu, praktiki olaraq heç vaxt deməkdir.
"X-Accel-Buffering": "no",
},
});
}
Müştəri tərəfi arxitekturası ardınca məlumatlar daxil olduqca dinamik, bloklamayan iz qrafiklərini render etmək üçün EventSource və React Suspense-dən istifadə edir. Bu, məlumat aliminə bütün məlumat bazasının brauzer yaddaşına yüklənməsini tələb etmədən zəncir qarışmasını və stasionarlığını dərhal vizual olaraq təsdiqləməyə imkan verir.
Yekunlaşdırsaq, müasir ehtimal konveyerinin mühəndisliyi ciddi sistem ayrılması (decoupling) məşğələsidir. Go-nun aşağı gecikməli asinxronluğunu, gRPC-nin binar səmərəliliyini, PostgreSQL pgvector-un qabaqcıl çoxölçülü indeksləşdirilməsini və Next.js Server-Sent Events-in qüsursuz axın imkanlarını birləşdirərək, qabaqcıl statistik modelləri ən ciddi istehsalat mühitlərinə yerləşdirməyə qadir olan davamlı bir arxitektura qururuq.
Mənbələr
27
- MCMC Methods: From Theory to Distributed Hamiltonian Monte Carlo over PySpark - Algorithms (MDPI)doi.org
- Rank-normalization, folding, and localization: An improved R-hat for assessing convergence of MCMC - Vehtari, Gelman, Simpson, Carpenter and Burknerarxiv.org
- Rank-normalization, folding, and localization: online appendix and codeavehtari.github.io
- posterior::rhat - Rhat convergence diagnostic, Stanmc-stan.org
- arviz.rhat - ArviZ documentationpython.arviz.org
- Threshold for R-hat (1.01 or 1.05) - stan-dev/rstan issue 812github.com
- Revisiting the Gelman-Rubin Diagnostic - arXivarxiv.org
- Gelman and Rubin Diagnostics - SAS/STAT user's guide extract, Imperial College ICICimperial.ac.uk
- Stopping Rules for Monte Carlo Methods: A Review - arXivarxiv.org
- High Performance gRPC - FOSDEM 2025archive.fosdem.org
- REST vs gRPC Performance in Go: A Practical Benchmark-Driven Guide - DEV Communitydev.to
- planetscale/vtprotobuf: A Protocol Buffers compiler that generates optimized codegithub.com
- gRPC Go: Basics tutorial and server-side streaminggrpc.io
- Go Documentation: the context packagepkg.go.dev
- Go Blog: Go Concurrency Patterns, pipelines and cancellationgo.dev
- cgo - Go Command Documentationpkg.go.dev
- Statically compiled Go programs, always, even with cgo, using musl - Dominik Honnefhonnef.co
- Multi-stage builds - Docker Docsdocs.docker.com
- pgvector/pgvector: Open-source vector similarity search for Postgresgithub.com
- An early look at HNSW performance with pgvector - Jonathan Katzjkatz05.com
- Writing R Extensions: linking against the R librarycran.r-project.org
- Guides: Streaming - Next.jsnextjs.org
- Route Handlers - Next.jsnextjs.org
- Real-Time Notifications with Server-Sent Events (SSE) in Next.js - Pedro Alonsopedroalonso.net
- MDN: Using server-sent eventsdeveloper.mozilla.org
- nginx: ngx_http_proxy_module, proxy_buffering and X-Accel-Bufferingnginx.org
- OpenTelemetry: distributed tracing conceptsopentelemetry.io



