Strömmande RPC:er och backpressure
Hantera server-, klient- och dubbelriktad strömning med @GrpcStreamMethod och observables.
Strömmande RPC:er och backpressure är en gratis lektion i NestJS: backend-API:er för företag på CoddyKit. Detta är lektion 3 av 4. Ni kan läsa hela lektionen gratis nedan och sedan öva praktiskt i webbläsaren med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt. Den ingår i lärvägen för NestJS: backend-API:er för företag, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i NestJS: backend-API:er för företag innehåller totalt 4 lektioner.
Fyra former av gRPC-anrop
gRPC definierar fyra typer av RPC, som skiljer sig åt beroende på om varje sida skickar ett enskilt meddelande eller en ström:
- Unärt — en begäran och ett svar. Standardfallet, som modelleras i NestJS med en vanlig metod som returnerar ett värde eller ett
Observable. - Serverströmning — en begäran, varefter servern skickar en ström av svar.
- Klientströmning — klienten skickar en ström av begäranden och servern svarar en gång.
- Tovegsströmning — båda sidor strömmar oberoende av varandra över samma HTTP/2-anslutning.
Strömning är viktigt för företags-API:er: direktsända prisflöden, loggströmning, uppladdning av filblock och chatt passar alla naturligt in i någon av de tre strömningsformerna i stället för polling.
Deklarera strömmar i Protobuf
Nyckelordet stream i tjänstedefinitionen i .proto väljer anropsformen. Placera det på begäran, svaret eller båda.
NestJS läser det här kontraktet vid uppstart och kopplar varje RPC till en hanterare. En metod vars svar är stream måste implementeras som en @GrpcStreamMethod (eller @GrpcStreamCall), inte som en vanlig @GrpcMethod.
syntax = "proto3";
package trading;
service PriceService {
// server streaming: one subscribe, many ticks
rpc Subscribe (SubscribeRequest) returns (stream PriceTick);
// client streaming: many orders, one ack
rpc PlaceOrders (stream Order) returns (OrderAck);
// bidirectional: stream in, stream out
rpc Chat (stream ChatMessage) returns (stream ChatMessage);
}
message SubscribeRequest { string symbol = 1; }
message PriceTick { string symbol = 1; double price = 2; int64 ts = 3; }
message Order { string id = 1; int32 qty = 2; }
message OrderAck { int32 accepted = 1; }
message ChatMessage { string user = 1; string text = 2; }Serverströmning med @GrpcMethod
Serverströmning returnerar många meddelanden för en begäran. I NestJS implementerar Ni detta som en vanlig @GrpcMethod som returnerar ett Observable. Varje värde som Observable-objektet skickar ut skickas som ett gRPC-meddelande; complete() stänger strömmen och error() avslutar den med en status.
Här skickar ett prisflöde ut en uppdatering varje sekund med RxJS interval och stoppar efter tio uppdateringar.
import { Controller } from '@nestjs/common';
import { GrpcMethod } from '@nestjs/microservices';
import { Observable, interval, map, take } from 'rxjs';
interface SubscribeRequest { symbol: string; }
interface PriceTick { symbol: string; price: number; ts: number; }
@Controller()
export class PriceController {
@GrpcMethod('PriceService', 'Subscribe')
subscribe(req: SubscribeRequest): Observable<PriceTick> {
return interval(1000).pipe(
take(10),
map((i) => ({
symbol: req.symbol,
price: 100 + Math.random(),
ts: Date.now() + i,
})),
);
}
}Klientströmning med @GrpcStreamMethod
När klienten strömmar skickar NestJS ett Observable med inkommande meddelanden till Er hanterare. Ni subscribe:ar på det, aggregerar meddelandena och löser ett enda svar när indataflödet är klart.
Det centrala mönstret är att returnera ett Promise (eller Subject) som Ni löser inuti complete. Använd @GrpcStreamMethod så att NestJS prenumererar på begärans Observable åt Er.
import { Controller } from '@nestjs/common';
import { GrpcStreamMethod } from '@nestjs/microservices';
import { Observable } from 'rxjs';
interface Order { id: string; qty: number; }
interface OrderAck { accepted: number; }
@Controller()
export class OrderController {
@GrpcStreamMethod('PriceService', 'PlaceOrders')
placeOrders(messages: Observable<Order>): Promise<OrderAck> {
return new Promise((resolve) => {
let accepted = 0;
messages.subscribe({
next: (order) => {
if (order.qty > 0) accepted++;
},
complete: () => resolve({ accepted }),
});
});
}
}Tovegsströmning
I tovegsströmning är båda Observable-flödena aktiva samtidigt. NestJS ger Er det inkommande Observable-objektet och förväntar sig att Ni returnerar ett utgående Observable-objekt (vanligtvis ett Subject som Ni matar med värden när inkommande meddelanden anländer).
Det här är formen för eko och chatt: prenumerera på det inkommande flödet, transformera meddelandena och anropa next() på det utgående Subject-objektet. Anropa complete() på det utgående Subject-objektet när det inkommande flödet är klart, så stängs det korrekt.
import { Controller } from '@nestjs/common';
import { GrpcStreamMethod } from '@nestjs/microservices';
import { Observable, Subject } from 'rxjs';
interface ChatMessage { user: string; text: string; }
@Controller()
export class ChatController {
@GrpcStreamMethod('PriceService', 'Chat')
chat(messages: Observable<ChatMessage>): Observable<ChatMessage> {
const out = new Subject<ChatMessage>();
messages.subscribe({
next: (m) =>
out.next({ user: 'server', text: `echo: ${m.text}` }),
complete: () => out.complete(),
error: (e) => out.error(e),
});
return out.asObservable();
}
}@GrpcStreamMethod kontra @GrpcStreamCall
NestJS erbjuder två dekoratorer för strömningshanterare:
- @GrpcStreamMethod — anpassad för RxJS. Ni tar emot ett
Observableoch returnerar ettObservable/Promise. NestJS hanterar det underliggande gRPC-anropsobjektet. - @GrpcStreamCall — på lägre nivå. Ni tar emot det råa gRPC-
call-objektet (en Node-duplexström) och styr det medcall.on('data'),call.write()ochcall.end().
Använd @GrpcStreamCall när Ni behöver direkt kontroll över strömmen — till exempel för att tillämpa backpressure via duplexströmmens returvärde från write(), vilket Observable-omslaget döljer.
Vad backpressure faktiskt innebär
Backpressure uppstår när en producent genererar data snabbare än konsumenten (eller nätverket) kan ta emot den. Utan hantering buffras överskottet i minnet tills processen får slut på heap-minne och kraschar.
gRPC över HTTP/2 har flödeskontroll: varje ström har ett fönster, och en långsam läsare slutar att utöka det. Node-duplexströmmen visar detta genom att write() returnerar false när den interna bufferten är full, samt genom en 'drain'-händelse när det är säkert att fortsätta.
RxJS-Observable-objekt är push-baserade och har ingen inbyggd backpressure — en interval fortsätter att skicka värden oavsett om socketen hinner med. Det är den centrala utmaningen i strömmande RPC-anrop.
Respektera backpressure från write()
Med @GrpcStreamCall får Ni den råa duplexströmmen och kan respektera flödeskontrollen korrekt. Regeln är: när call.write(msg) returnerar false ska Ni sluta producera tills händelsen 'drain' inträffar.
Den här generatorstyrda pumpen hämtar nästa objekt först efter att den föregående skrivningen har accepterats, så minnesanvändningen förblir begränsad oavsett hur snabbt källan kan producera data.
import { Controller } from '@nestjs/common';
import { GrpcStreamCall } from '@nestjs/microservices';
@Controller()
export class FeedController {
@GrpcStreamCall('PriceService', 'Subscribe')
subscribe(call: any) {
let i = 0;
const pump = () => {
let ok = true;
while (ok && i < 100000) {
const tick = { symbol: 'ACME', price: i, ts: Date.now() };
ok = call.write(tick); // false => buffer full
i++;
}
if (i < 100000) {
call.once('drain', pump); // resume when flushed
} else {
call.end();
}
};
pump();
}
}RxJS-operatorer som hanterar backpressure
När Ni stannar i Observable-världen kan Ni inte pausa en het källa, men Ni kan forma genomströmningen så att konsumenten inte överbelastas:
concatMap— bearbetar en inre uppgift helt innan nästa påbörjas; bevarar ordningen och serialiserar arbetet.throttleTime/sampleTime— slänger mellanliggande värden och skickar högst ett värde per tidsfönster (bra för brusiga prisflöden).bufferTime/bufferCount— samlar många utskick till ett meddelande och minskar kostnaden per meddelande.auditTime— skickar det senaste värdet efter ett lugnt tidsfönster.
Dessa operatorer byter fullständighet eller fördröjning mot en begränsad belastning. mergeMap utan en gräns för samtidighet gör motsatsen — den skapar ett obegränsat antal parallella operationer och är en vanlig fallgrop när det gäller backpressure.
import { interval, Observable } from 'rxjs';
import { sampleTime, map, take } from 'rxjs';
interface PriceTick { symbol: string; price: number; }
// Raw feed could emit every 1ms; downstream only needs ~10/sec.
function throttledFeed(): Observable<PriceTick> {
return interval(1).pipe(
map((i) => ({ symbol: 'ACME', price: 100 + (i % 50) })),
sampleTime(100), // at most one tick per 100ms
take(20),
);
}Begränsa samtidigheten med concatMap
Ett vanligt misstag i företagssystem är att mappa varje objekt i ett inkommande flöde till ett asynkront anrop med mergeMap och obegränsad samtidighet. Under belastning samlas tusentals pågående promises och tömmer DB-poolen.
Föredra concatMap (samtidighet 1) eller mergeMap(fn, N) med en uttrycklig gräns. Kodavsnittet nedan bearbetar inkommande beställningar strikt en i taget, vilket ger naturlig backpressure mot persistenslagret.
import { Observable, from } from 'rxjs';
import { concatMap, toArray } from 'rxjs';
interface Order { id: string; qty: number; }
// Simulate a slow persistence call.
function persist(o: Order): Promise<string> {
return new Promise((res) => setTimeout(() => res(o.id), 5));
}
function processOrders(orders: Observable<Order>): Observable<string[]> {
return orders.pipe(
concatMap((o) => from(persist(o))), // one at a time
toArray(),
);
}Avbrytning, tidsgränser och städning
Strömmar måste avslutas korrekt även när klienten kopplar från eller en tidsgräns löper ut. Om Ni ignorerar avbrytning fortsätter producenten att skicka data till en död socket och läcker timers och DB-pekare.
- Med
@GrpcStreamCalllyssnar Ni eftercall.on('cancelled', ...)och stoppar pumpen. - Med Observable-objekt körs RxJS-nedmonteringen när NestJS avslutar prenumerationen vid avbrytning — placera städningen i en
finalize-operator eller i producentens nedmonteringsfunktion. - Vidarebefordra alltid fel via
error()(eller kasta enRpcException) så att klienten ser en riktig gRPC-status i stället för en ström som hänger tyst.
import { Observable, interval } from 'rxjs';
import { map, takeWhile, finalize } from 'rxjs';
interface Tick { price: number; }
function feedWithCleanup(): Observable<Tick> {
return interval(500).pipe(
map((i) => ({ price: 100 + i })),
takeWhile((t) => t.price < 110),
finalize(() => {
// runs on complete, error, OR client cancel/unsubscribe
console.log('feed torn down, releasing resources');
}),
);
}Snabbkontroll: välj dekorator
Ni strömmar en export med flera miljoner rader till en gRPC-klient och ser hur Node-processens heap växer tills den får slut på minne. Ni misstänker att producenten kör ifrån socketen.
Sammanfattning
Nu vet Ni hur man implementerar och styr strömmande RPC-anrop i NestJS:
- Nyckelordet
streami.protoväljer serverströmning, klientströmning eller tovegsströmning. @GrpcMethodsom returnerar ettObservablehanterar serverströmning;@GrpcStreamMethodger Er ett inkommandeObservableför klient- och tovegsströmning (returnera ettPromiseellerSubject).- RxJS är push-baserat och har ingen inbyggd backpressure — operatorer som
concatMap,sampleTimeochbufferTimeformar belastningen men pausar inte en het källa. - För riktig flödeskontroll går Ni ner till
@GrpcStreamCalloch respekterar returvärdetfalsefrånwrite()samt händelsen'drain'. - Hantera avbrytning och tidsgränser med
cancelled-lyssnare ellerfinalize, och exponera fel viaerror()/RpcExceptionså att klienterna får en riktig status.
Lär dig TypeScript med en AI-lärare – gratis
Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.
- Kurser
- 20
- Lektioner
- 76
Vanliga frågor
Är lektionen ”Strömmande RPC:er och backpressure” gratis?
Ja – hela texten till ”Strömmande RPC:er och backpressure” kan läsas gratis här på webben. Om Ni vill öva interaktivt med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt och låsa upp resten av kursen i NestJS: backend-API:er för företag, kan Ni uppgradera till CoddyKit PRO. Kursen i NestJS: backend-API:er för företag innehåller totalt 4 lektioner.
Vad lär jag mig i ”Strömmande RPC:er och backpressure”?
Hantera server-, klient- och dubbelriktad strömning med @GrpcStreamMethod och observables. Ni övar på NestJS: backend-API:er för företag med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.
Behöver jag någon erfarenhet för att börja lära mig NestJS: backend-API:er för företag?
Du behöver inga förkunskaper. Utbildningen i NestJS: backend-API:er för företag på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 3 av 4.
Hur lång tid tar lektionen ”Strömmande RPC:er och backpressure”?
De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.
Kan jag skriva och köra kod i den här NestJS: backend-API:er för företag-lektionen?
Ja. Varje NestJS: backend-API:er för företag-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.
Alla lektioner i den här kursen
- Definiera tjänster och meddelanden i Protobuf
- Implementera och använda gRPC-metoder
- Strömmande RPC:er och backpressure
- Utveckling av kontrakt och bakåtkompatibilitet