Streaming-RPC's en backpressure
Verwerk server-, client- en bidirectionele streaming met @GrpcStreamMethod en observables.
Streaming-RPC's en backpressure is een gratis Enterprise-backend-API's met NestJS-les op CoddyKit. Dit is les 3 van 4. Je kunt de volledige les hieronder gratis lezen en daarna in de browser praktisch oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject Enterprise-backend-API's met NestJS. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus Enterprise-backend-API's met NestJS bevat in totaal 4 lessen.
Vier vormen van gRPC-aanroepen
gRPC definieert vier soorten RPC's, afhankelijk van de vraag of elke kant één enkel bericht of een stream verstuurt:
- Uniair — één verzoek, één antwoord. De standaardvorm; in NestJS gemodelleerd met een gewone methode die een waarde of
Observableteruggeeft. - Serverstreaming — één verzoek, waarna de server een stream met antwoorden verstuurt.
- Clientstreaming — de client verstuurt een stream met verzoeken, waarna de server één keer antwoordt.
- Bidirectioneel — beide kanten streamen onafhankelijk via dezelfde HTTP/2-verbinding.
Streaming is belangrijk voor API's van bedrijfsdiensten: live prijsfeeds, logbestanden volgen, bestanden in delen uploaden en chatten passen allemaal vanzelf in een van de drie streamingvormen, in plaats van dat je moet pollen.
Streams declareren in Protobuf
Het sleutelwoord stream in de servicedefinitie van je .proto bepaalt de vorm van de aanroep. Plaats het bij het verzoek, het antwoord of beide.
NestJS leest dit contract bij het opstarten en koppelt elke RPC aan een handler. Een methode waarvan het antwoord stream is, moet worden geïmplementeerd als een @GrpcStreamMethod (of @GrpcStreamCall), niet als een gewone @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; }Serverstreaming met @GrpcMethod
Bij serverstreaming worden voor één verzoek meerdere berichten teruggegeven. In NestJS implementeer je dit als een normale @GrpcMethod die een Observable teruggeeft. Elke waarde die de Observable uitstuurt, wordt als één gRPC-bericht verzonden; complete() sluit de stream en error() beëindigt deze met een status.
Hier stuurt een prijsfeed elke seconde een meetpunt met RxJS interval en stopt deze na tien meetpunten.
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,
})),
);
}
}Clientstreaming met @GrpcStreamMethod
Wanneer de client streamt, geeft NestJS je handler een Observable met binnenkomende berichten. Je abonneert je erop met subscribe, verzamelt de berichten en lost één antwoord op zodra de invoerstroom is voltooid.
Het belangrijkste patroon: geef een Promise (of Subject) terug die je binnen complete oplost. Gebruik @GrpcStreamMethod, zodat NestJS de Observable met verzoeken voor je abonneert.
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 }),
});
});
}
}Bidirectionele streaming
In de bidirectionele modus zijn beide Observables tegelijk actief. NestJS geeft je de binnenkomende Observable en verwacht dat je een uitgaande Observable teruggeeft (meestal een Subject waar je waarden in zet zodra berichten binnenkomen).
Dit is de vorm voor echo en chat: abonneer je op de inkomende stream, transformeer de berichten en roep next() aan op het uitgaande subject. Roep complete() aan op het uitgaande subject wanneer de inkomende stream is voltooid, zodat deze netjes wordt afgesloten.
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 versus @GrpcStreamCall
NestJS biedt twee decorators voor streaminghandlers:
- @GrpcStreamMethod — geschikt voor RxJS. Je ontvangt een
Observableen geeft eenObservable/Promiseterug. NestJS beheert het onderliggende gRPC-aanroepobject. - @GrpcStreamCall — lager niveau. Je ontvangt de onbewerkte gRPC-
call(een Node-duplexstream) en bestuurt deze metcall.on('data'),call.write()encall.end().
Gebruik @GrpcStreamCall wanneer je directe controle over de stream nodig hebt, bijvoorbeeld om backpressure toe te passen via de geretourneerde waarde van write() op de duplexstream, die de Observable-wrapper voor je verbergt.
Wat backpressure werkelijk is
Backpressure ontstaat wanneer een producent sneller gegevens genereert dan de consument (of het netwerk) deze kan accepteren. Zonder afhandeling wordt het overschot in het geheugen gebufferd totdat het proces geen heapruimte meer heeft en crasht.
gRPC via HTTP/2 heeft flowcontrol: elke stream heeft een venster en een trage lezer zorgt ervoor dat dit venster niet verder opschuift. De Node-duplexstream maakt dit zichtbaar doordat write() false teruggeeft wanneer de interne buffer vol is, en een gebeurtenis 'drain' afgeeft wanneer het veilig is om verder te gaan.
RxJS-Observables werken volgens het pushmodel en hebben geen ingebouwde backpressure — een interval blijft waarden uitsturen, ongeacht of de socket kan bijhouden. Dat is de belangrijkste spanning bij streaming-RPC's.
Backpressure van write() respecteren
Met @GrpcStreamCall krijg je de onbewerkte duplexstream en kun je flowcontrol correct naleven. De regel is: wanneer call.write(msg) false teruggeeft, stop je met produceren totdat de gebeurtenis 'drain' optreedt.
Deze generatorgestuurde verwerker haalt het volgende item pas op nadat de vorige schrijfactie is geaccepteerd. Zo blijft het geheugengebruik begrensd, ongeacht hoe snel de bron zou kunnen produceren.
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-operators die backpressure ondersteunen
Wanneer je binnen de Observable-wereld blijft, kun je een actieve bron niet pauzeren, maar je kunt de doorvoer wel vormgeven, zodat de consument niet wordt overbelast:
concatMap— verwerkt elke interne taak volledig voordat de volgende begint; behoudt de volgorde en voert het werk na elkaar uit.throttleTime/sampleTime— laat tussenliggende waarden vallen en stuurt maximaal één waarde per venster uit (handig voor rumoerige prijsfeeds).bufferTime/bufferCount— bundelt veel waarden in één bericht, waardoor de overhead per bericht afneemt.auditTime— stuurt de meest recente waarde uit na een stil venster.
Deze operators ruilen volledigheid of latentie in voor een begrensde belasting. mergeMap zonder limiet voor gelijktijdigheid doet het tegenovergestelde: het verspreidt werk onbeperkt en is een veelvoorkomende valkuil bij 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),
);
}Begrensde gelijktijdigheid met concatMap
Een veelgemaakte fout in bedrijfssoftware is elk binnenkomend item uit een stream met mergeMap aan een asynchrone aanroep te koppelen, met onbeperkte gelijktijdigheid. Bij belasting stapelen duizenden actieve promises zich op en raakt de databasepool uitgeput.
Gebruik liever concatMap (gelijktijdigheid 1) of mergeMap(fn, N) met een expliciete limiet. Het onderstaande fragment verwerkt binnenkomende orders strikt één voor één, waardoor er vanzelf backpressure ontstaat op de persistentielaag.
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(),
);
}Annulering, deadlines en opruimen
Streams moeten netjes eindigen, ook wanneer de client de verbinding verbreekt of een deadline verstrijkt. Als je annulering negeert, blijft je producent gegevens naar een dode socket sturen en lekken timers en databasecursors.
- Luister met
@GrpcStreamCallnaarcall.on('cancelled', ...)en stop je verwerker. - Bij Observables wordt de opruimactie van RxJS uitgevoerd wanneer NestJS bij annulering de abonnementen opzegt — plaats de opruimactie in een
finalize-operator of in de opruimfunctie van de producent. - Geef fouten altijd door via
error()(of gooi eenRpcException), zodat de client een echte gRPC-status ziet in plaats van een stil vastgelopen stream.
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');
}),
);
}Snelle controle: de decorator kiezen
Je streamt een export met miljoenen rijen naar een gRPC-client en ziet dat de heap van het Node-proces blijft groeien totdat er een out-of-memory-fout optreedt. Je vermoedt dat de producent sneller is dan de socket kan verwerken.
Samenvatting
Je weet nu hoe je streaming-RPC's in NestJS implementeert en beheert:
- Het sleutelwoord
streamin.protoselecteert server-, client- of bidirectionele streaming. @GrpcMethoddie eenObservableteruggeeft, verwerkt serverstreaming;@GrpcStreamMethodgeeft je een binnenkomendeObservablevoor client- en bidirectionele streams (geef eenPromiseofSubjectterug).- RxJS werkt volgens het pushmodel en heeft geen ingebouwde backpressure — operators zoals
concatMap,sampleTimeenbufferTimeregelen de belasting, maar pauzeren geen actieve bron. - Voor echte flowcontrol ga je naar
@GrpcStreamCallen respecteer je de retourwaardefalsevanwrite()en de gebeurtenis'drain'. - Verwerk annulering en deadlines met listeners voor
cancelledoffinalize, en geef fouten door viaerror()/RpcException, zodat clients een echte status ontvangen.
Leer TypeScript met een AI-tutor — gratis
Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.
- Cursussen
- 20
- Lessen
- 76
Veelgestelde vragen
Is de les “Streaming-RPC's en backpressure” gratis?
Ja — de volledige tekst van “Streaming-RPC's en backpressure” kun je hier gratis op het web lezen. Als je interactief wilt oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is, en de rest van de cursus Enterprise-backend-API's met NestJS wilt ontgrendelen, kun je upgraden naar CoddyKit PRO. De cursus Enterprise-backend-API's met NestJS bevat in totaal 4 lessen.
Wat leer ik in “Streaming-RPC's en backpressure”?
Verwerk server-, client- en bidirectionele streaming met @GrpcStreamMethod en observables. Je oefent met Enterprise-backend-API's met NestJS door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.
Heb ik ervaring nodig om met Enterprise-backend-API's met NestJS te beginnen?
Ervaring vooraf is niet nodig. Enterprise-backend-API's met NestJS op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 3 van 4.
Hoe lang duurt de les “Streaming-RPC's en backpressure”?
De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.
Kan ik code schrijven en uitvoeren in deze les over Enterprise-backend-API's met NestJS?
Ja. Elke les over Enterprise-backend-API's met NestJS bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.
Alle lessen in deze cursus
- Services en berichten definiëren in Protobuf
- gRPC-methoden implementeren en gebruiken
- Streaming-RPC's en backpressure
- Contractevolutie en achterwaartse compatibiliteit