Strona główna / Artykuły / Zrozumienie strumieni w Node.js: problem, który faktycznie rozwiązują

Zrozumienie strumieni w Node.js: problem, który faktycznie rozwiązują

Dowiedz się, dlaczego istnieją strumienie w Node.js, jak funkcjonuje przekierowanie danych wewnątrz systemu oraz co tak naprawdę oznacza ciśnienie zwrotne przy efektywnym obsługiwaniu dużych ilości danych.

1261 słów

Większość poradników na temat strumieni zaczyna się od omówienia interfejsu API, takiego jak .pipe(), Readable, Writable, Transform, zanim w ogóle zajrzy się w problem, który te narzędzia mają rozwiązać. Taki sposób podejścia sprawia, że strumienie wydają się arbitralnymi procedurami, które trzeba zapamiętać. Odwróćmy kolejność – zacznijmy od problemu, a interfejs API w większości przypadków sam się wyjaśni.

Problem: Niektóre dane są zbyt duże, by przechowywać je jednocześnie

Załóżmy zadanie polegające na przekształceniu pliku CSV o rozmiarze 4 GB w format JSON. Proste podejście wygląda tak:

const fs = require("fs");

const data = fs.readFileSync("export.csv", "utf8");
const rows = data.split("\n").map(parseRow);
fs.writeFileSync("export.json", JSON.stringify(rows));

readFileSync nie zwróci kontroli programowi, dopóki cały plik o wielkości 4 GB nie zostanie załadowany do pamięci jako pojedynczy ciąg znaków w JavaScript, a taka reprezentacja ciągu zazwyczaj wymaga znacznie więcej pamięci, niż sugerowałaby sama wielkość pliku w formacie surowym. Na maszynie z zaledwie 2 GB dostępnej pamięci RAM ten kod nie tylko działa nieskutecznie, ale także powoduje awarię lub blokuje wszystkie inne procesy walczące o tę samą pulę pamięci. W samej logice nie ma nic złego; wszystkie kroki transformacji są poprawne. Prawdziwym błędem jest ukryte założenie wewnątrz readFileSync: że można jednocześnie przechowywać cały plik w pamięci, bez względu na to, jak duży jest ten plik.

Prawdziwa idea za strumieniem

Zamiast żądać całego zestawu danych od razu, strumieniowanie opiera się na innej zasadzie: żąda się jednego fragmentu, przetwarza go, a następnie kolejnego. W żadnym momencie pełny plik nie znajduje się w pamięci. W danym momencie istnieje tylko mały jego fragment, który jest przetwarzany i uwalniany, zanim pojawi się następny.

const fs = require("fs");

const readStream = fs.createReadStream("export.csv", { encoding: "utf8" });

readStream.on("data", (chunk) => {
  console.log(`Received ${chunk.length} characters`);
});

readStream.on("end", () => {
  console.log("Done reading the whole file, piece by piece");
});

Zauważ, że całkowita wielkość pliku w ogóle nie odgrywa roli w tym kodzie. Niezależnie od tego, czy plik źródłowy ma 4 GB, czy 4 KB, ten sam fragment kodu działa identycznie, chociaż maksymalne zużycie pamięci różni się ogromnie w tych dwóch przypadkach. W tym właśnie polega sens strumieniowania: zamiast wymagać posiadania całego zestawu danych przed rozpoczęciem pracy, umożliwia on natychmiastowe rozpoczęcie i zakończenie zadania bez przechowywania w pamięci więcej niż tylko niewielkiego zakresu danych.

Piping: Połączenie źródła bezpośrednio z docelowym miejscem

Ręczne pobieranie fragmentów po jednym jest przydatną techniką, ale w praktyce znacznie częściej łączy się strumień do odczytu bezpośrednio z tym do zapisu, umożliwiając automatyczny przepływ danych od źródła do celu zamiast ręcznego przekazywania każdego fragmentu:

const fs = require("fs");

fs.createReadStream("export.csv")
  .pipe(fs.createWriteStream("export-copy.csv"));

Dwie linijki kodu wystarczą do skopiowania pliku dowolnej wielkości bez konieczności ładowania całego pliku do pamięci. Funkcja .pipe() nie robi tu żadnych „czarów” – po prostu kieruje zdarzenie data emitowane przez źródło do funkcji write w celu, a także zajmuje się dodatkowym detalami, które okazują się ważniejsze niż sam mechanizm kopiowania.

Część, której prawie nikt nie wyjaśnia dobrze: ciśnienie wsteczne

To jest prawdziwy problem, który rozwiązuje .pipe(), a nie tylko kwestia wygody. Wyobraź sobie sytuację, w której odczyt odbywa się szybko, na przykład z lokalnego dysku, podczas gdy zapis jest wolny, być może przez połączenie sieciowe o ograniczonej przepustowości.

readStream.on("data", (chunk) => {
  writeStream.write(chunk); // what happens if this can't keep up?
});

Wywoływanie .write() szybciej, niż destynacja jest w stanie przetworzyć dane, nie powoduje, że strumień do zapisu rzuci błąd ani zablokuje wykonywanie. Zamiast tego w tajemnicy gromadzi nadmiarowe dane w wewnętrznym buforze pamięci, przechowując je aż do chwili, gdy będzie mógł je wyrzucić. Jeśli ta niespójność trwa wystarczająco długo, przy dużej różnicy między szybkością odczytu a zapisu, dokładnie ten problem, który miało rozwiązać streamowanie, powraca w ukryty sposób: zużycie pamięci rośnie bez ograniczeń, choć jest odroczone, a nie natychmiastowe, i zazwyczaj znacznie trudniej jest go zauważyć, zanim doprowadzi do awarii.

Backpressure to zabezpieczenie przed właśnie tym modelem awarii: daje strumieniowi do zapisu możliwość sygnalizowania, że osiągnął swoją pojemność i wymaga od producenta zmniejszenia tempa nadawania danych, a prawidłowo zaimplementowany producent uwzględnia tę sygnalizację zamiast nadal przesyłać dane.

readStream.on("data", (chunk) => {
  const canContinue = writeStream.write(chunk);
  if (!canContinue) {
    readStream.pause(); // stop reading until the writable side catches up
  }
});

writeStream.on("drain", () => {
  readStream.resume(); // writable side is ready for more
});

Gdy tylko wewnętrzna bufora strumienia do zapisu przekroczy ustawiony limit, metoda .write() zwraca false, co służy jako sygnał do wstrzymania nadawania danych, dopóki strumień nie wygeneruje zdarzenia drain, wskazującego, że usunął zaległości i jest gotowy przyjąć więcej danych. To właśnie ten wzorzec pauzowania i kontynuowania pracy jest automatycznie obsługiwany przez metodę .pipe() w tle:

readStream.pipe(writeStream); // handles backpressure for you, silently, correctly

To, a nie zwięzłość, jest prawdziwym uzasadnieniem faworyzowania .pipe() zamiast ręcznego łączenia słuchaczy data i write. Przekierowywanie fragmentów danych ręcznie, ignorując wartość zwracaną przez .write(), ponownie powoduje ten sam problem nieograniczonej pamięci, którego miały uniknąć strumienie – jest to po prostu błąd bliższy temu, który istnieje w readFileSync.

Strumienie transformacyjne: przetwarzanie danych w trakcie przesyłania

Są sytuacje, gdy samo przeniesienie danych z jednego miejsca do drugiego bez żadnych zmian nie wystarcza – trzeba je jeszcze w drodze przekształcić. Strumień Transform został stworzony właśnie do tego: znajduje się pośrodku łańcucha strumieni, przyjmuje napływające fragmenty danych, aplikuje do każdego z nich określoną operację i przekazuje wynik dalej:

const { Transform } = require("stream");

const upperCaseTransform = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  },
});

fs.createReadStream("input.txt")
  .pipe(upperCaseTransform)
  .pipe(fs.createWriteStream("output.txt"));

Każdy fragment jest przekształcany na wielkie litery w trakcie przechodzenia przez łańcuch przetwarzania, a w żadnym momencie cały plik, ani w swojej oryginalnej formie, ani po przekształceniu, nie musi być w całości przechowywany w pamięci. To właśnie jest mechanizm stojący za wbudowanym w Node funkcją zlib.createGzip(). Jest to nic innego jak strumień Transform, który kompresuje każdy przychodzący fragment, a podobnie jak każda inna transformacja, może być bezpośrednio włączony do łańcucha przetwarzania w taki sam sposób jak przykład z wielkimi literami:

const zlib = require("zlib");

fs.createReadStream("export.csv")
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream("export.csv.gz"));

Czytanie, kompresowanie i zapiszanie odbywają się jednocześnie, przetwarzając małe fragmenty danych w miarę ich przychodzenia, przy czym żdy z tych trzech etapów nie wymaga jednoczesnego załadowania całego pliku.

Dlaczego warto to faktycznie zrozumieć

Strumienie są uważane za jeden z najtrudniejszych obszarów interfejsu API Node, a szczerze mówiąc, surowy interfejs oparty na zdarzeniach rzeczywiście wydaje się niezgrabny, więc ta reputacja nie jest całkowicie bezpodstawna. Jednak koncepcja leżąca u jego podstaw jest prosta: unikać ładowania całego zbioru danych do pamięci, przetwarzać je kawałek po kawałku i upewniać się, że szybki producent nie może w ten sposób cicho zalewać wolnego konsumenta. Gdy już przyjmiesz to jako właściwy model, .pipe(), Transform oraz mechanizm backpressure przestaną wydawać się trzema niepowiązanymi interfejsami do osobnego zapamiętywania. W rzeczywistości stanowią one jedną ideę, prezentowaną poprzez trzy powiązane elementy API, rozwiązując problem, z którym coś takiego jak readFileSync nigdy nie musiało się mierzyć, ponieważ od początku nie było przeznaczone do obsługi niczego poza małymi plikami.

Literatura pokrewna