İçeriğe geç
academia.sh

Ders 07 / 20

Akışlar

Dört akış türü, parça sınırlarının kayıt sınırlarıyla örtüşmemesi, nesne kipi, boru hattı kurma ve geri basıncın ölçülmesi.

İçindekiler

Önceki iki ders dosyayı bir kerede belleğe aldı. On iki satırlık bir dosyada bunun bedeli yoktur. Ölçüm düğümleri günlerce yazdığında dosya yüzlerce megabayta ulaşır ve aynı kod, dosyanın tamamı kadar bellek ister — sonra o belleği bir kez daha kopyalayarak satırlara böler.

Bu dersin sorusu şudur: veriyi tamamını görmeden nasıl işleriz? Yanıt, Kabuk Programlama kursundaki boru hattı modelinin aynısıdır — veri bir uçtan girer, birkaç dönüşümden geçer, öbür uçtan çıkar; hiçbir adım tamamını tutmaz. Çalışma zamanı bu modeli birinci sınıf nesnelerle sunar.

Dört Tür

Akış (stream), veriyi parça parça taşıyan bir soyutlamadır. Dört türü vardır:

  • Okunabilir akış, veri üretir. Dosya okuma, gelen HTTP isteği, alt sürecin standart çıktısı bu türdendir.
  • Yazılabilir akış, veri tüketir. Dosya yazma, giden HTTP yanıtı, standart çıktı bu türdendir.
  • Çift yönlü akış, ikisi birden olan bağımsız iki yön taşır; ağ soketi böyledir.
  • Dönüştürücü akış, çift yönlünün özel hâlidir: yazılan veriyi işleyip okunabilir tarafına çıkarır. Sıkıştırma, şifreleme ve ayrıştırma bu türdendir.

Bütün akışlar olay yayıcıdır: data, end, error, drain gibi olaylar yayarlar. Olay yayıcının kendisi sonraki derslerden birinin konusudur; burada olayları doğrudan dinlemek yerine daha üst düzey araçlar kullanılacak.

Parça Sınırı Kayıt Sınırı Değildir

Okunabilir akış, veriyi kendi belirlediği büyüklükte parçalar hâlinde verir. Bu büyüklük highWaterMark seçeneğiyle ayarlanır ve dosyadaki satır sınırlarıyla hiçbir ilgisi yoktur.

// akis-parca.mjs
import { createReadStream } from 'node:fs';

const akis = createReadStream('olcumler.ndjson', { highWaterMark: 256 });

let sira = 0;
for await (const parca of akis) {
  sira += 1;
  console.log(`parca ${sira}: ${parca.length} bayt, tur ${parca.constructor.name}`);
}
console.log(`toplam ${sira} parca`);
node akis-parca.mjs
parca 1: 256 bayt, tur Buffer
parca 2: 256 bayt, tur Buffer
parca 3: 256 bayt, tur Buffer
parca 4: 232 bayt, tur Buffer
toplam 4 parca

Dosya 1000 bayttır; dört parçaya bölündü ve son parça kalanı taşıdı. Parçaların hiçbiri bir satır sınırında bitmiyor. Bu, akışla çalışmanın temel zorluğudur: kayıt sınırını kurmak akışın değil, sizin işinizdir.

Örnekteki for await döngüsü, okunabilir akışın eşzamansız yinelenebilir olmasından yararlanır. Döngü gövdesi çalışırken akış duraklar; bu, geri basıncın en yalın uygulamasıdır.

Nesne Kipi ve Dönüştürücü

Varsayılan olarak akışlar bayt taşır. Nesne kipi (object mode) açıldığında parça yerine keyfi bir değer taşınır; böylece boru hattının bir yerinden sonra baytlar değil kayıtlar akar.

Aşağıdaki dönüştürücü, bayt parçalarını satırlara böler ve yarım kalan satırı bir sonraki parçayla birleştirmek üzere kendinde tutar. Bu, akış yazımının klasik kalıbıdır.

// akis-ozet.mjs
import { createReadStream } from 'node:fs';
import { Transform, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';

// Bayt parcalarini satirlara boler; yarim kalan satiri kendinde tutar
class SatirBolucu extends Transform {
  #kalan = '';

  constructor() {
    super({ readableObjectMode: true });
  }

  _transform(parca, kodlama, bitti) {
    const metin = this.#kalan + parca.toString('utf8');
    const satirlar = metin.split('\n');
    this.#kalan = satirlar.pop();          // son parca yarim olabilir
    for (const satir of satirlar) {
      if (satir !== '') this.push(satir);
    }
    bitti();
  }

  _flush(bitti) {
    if (this.#kalan !== '') this.push(this.#kalan);
    bitti();
  }
}

// Satirlari kayda cevirip kova kova toplar
class OzetToplayici extends Writable {
  kovalar = new Map();

  constructor() {
    super({ objectMode: true });
  }

  _write(satir, kodlama, bitti) {
    let kayit;
    try {
      kayit = JSON.parse(satir);
    } catch (hata) {
      return bitti(new Error(`bozuk satir: ${satir.slice(0, 30)}`, { cause: hata }));
    }
    const anahtar = `${kayit.dugum}/${kayit.metrik}`;
    const kova = this.kovalar.get(anahtar) ?? { sayi: 0, toplam: 0 };
    kova.sayi += 1;
    kova.toplam += kayit.deger;
    this.kovalar.set(anahtar, kova);
    bitti();
  }
}

const toplayici = new OzetToplayici();
await pipeline(
  createReadStream('olcumler.ndjson', { highWaterMark: 256 }),
  new SatirBolucu(),
  toplayici,
);

for (const [anahtar, { sayi, toplam }] of toplayici.kovalar) {
  console.log(`${anahtar.padEnd(18)} ${String(sayi).padStart(2)} olcum, ortalama ${(toplam / sayi).toFixed(2)}`);
}
node akis-ozet.mjs
kenar-01/sicaklik   3 olcum, ortalama 21.87
kenar-01/nem        2 olcum, ortalama 47.60
kenar-02/sicaklik   3 olcum, ortalama 20.23
kenar-02/nem        1 olcum, ortalama 52.50
kenar-03/sicaklik   2 olcum, ortalama 24.35
kenar-03/nem        1 olcum, ortalama 41.30

Üç ayrıntı dikkat ister.

readableObjectMode yalnızca okunabilir tarafı nesne kipine alır; yazılabilir taraf bayt kabul etmeye devam eder. Dönüştürücünün iki tarafı ayrı yapılandırılabilir ve tam olarak bu örnek için gereklidir.

_flush, girdi bittiğinde bir kez çağrılır ve elde kalan yarım satırı çıkarma fırsatı verir. Bu kanca olmadan, son satırı satır sonu ile bitmeyen dosyaların son kaydı kaybolur.

bitti geri çağrısına bir hata verildiğinde, hata akışa yayılır ve boru hattını düşürür. Yazılabilir akış içinde throw kullanmak aynı sonucu vermez.

Boru Hattı Kurmak

Akışları birbirine bağlamanın iki yolu vardır. pipe yöntemi eskidir ve bir kusuru vardır: zincirin ortasındaki bir akış hata verdiğinde diğerleri kendiliğinden kapanmaz, açık dosya tanıtıcıları ve bellekte tutulan tamponlar geride kalır.

pipeline bu kusuru kapatır. Zincirdeki herhangi bir akış hata verdiğinde tümünü yok eder ve hatayı tek bir yerden bildirir.

Bozuk bir satır içeren dosyayla sınayalım. Önce dosyayı oluşturun:

printf '{"dugum":"kenar-01","metrik":"nem","deger":48.0}\nbozuk satir\n{"dugum":"kenar-02","metrik":"nem","deger":50.0}\n' > bozuk.ndjson
// akis-hata.mjs
import { createReadStream } from 'node:fs';
import { Transform, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';

const ayristir = new Transform({
  readableObjectMode: true,
  transform(parca, kodlama, bitti) {
    for (const satir of parca.toString('utf8').split('\n')) {
      if (satir === '') continue;
      try {
        this.push(JSON.parse(satir));
      } catch (hata) {
        return bitti(new Error(`ayristirilamayan satir: ${satir}`, { cause: hata }));
      }
    }
    bitti();
  },
});

let sayac = 0;
const say = new Writable({
  objectMode: true,
  write(kayit, kodlama, bitti) { sayac += 1; bitti(); },
});

const kaynak = createReadStream('bozuk.ndjson');

try {
  await pipeline(kaynak, ayristir, say);
  console.log('tamamlandi, kayit sayisi:', sayac);
} catch (hata) {
  console.log('boru hatti hatasi :', hata.message);
  console.log('asil neden        :', hata.cause.name);
  console.log('kaynak yok edildi :', kaynak.destroyed);
  console.log('hedef yok edildi  :', say.destroyed);
}
node akis-hata.mjs
boru hatti hatasi : ayristirilamayan satir: bozuk satir
asil neden        : SyntaxError
kaynak yok edildi : true
hedef yok edildi  : true

Son iki satır önemlidir: hata bildirildiğinde kaynak ve hedef zaten kapatılmıştı. Dosya tanıtıcısı sızmadı. node:stream/promises modülündeki söz döndüren biçim, bu temizliği try/catch ile birleştirir.

Geri Basınç

Veri Yapıları kursunda kuyruklar ele alınırken geri basınç (backpressure) tanımlanmıştı: üretici tüketiciden hızlıysa, aradaki tamponun büyümesini durdurmanın tek yolu üreticiyi yavaşlatmaktır. Aynı sorun Kabuk Programlama kursunda boru hattı anlatılırken de geçti — orada yavaşlatmayı çekirdek yapıyordu.

Akışlarda bu iş yazma çağrısının dönüş değerine emanet edilir. write çağrısı, tampondaki veri eşiği aştıysa false döner. Bu bir hata değildir, “yavaşla” demektir; üretici drain olayını bekleyip devam eder.

Kurs boyunca iki sözcük ayrı tutulacak. Tampon, bir akışın içinde işlenmeyi bekleyen veriyi tutan alandır; büyüklüğü writableLength alanından okunur. Arabellek ise ham baytları temsil eden nesnedir ve sonraki dersin konusudur. İngilizce her ikisi için de buffer sözcüğünü kullanır; ayrım Türkçe karşılıklarla korunuyor.

// geri-basinc.mjs
import { Writable } from 'node:stream';
import { once } from 'node:events';

// Her kaydi 5 ms'de isleyen, en fazla 4 kayit tamponlayan hedef
class YavasYazici extends Writable {
  constructor() {
    super({ objectMode: true, highWaterMark: 4 });
  }
  _write(kayit, kodlama, bitti) {
    setTimeout(bitti, 5);
  }
}

function kayitlar(adet) {
  return Array.from({ length: adet }, (_, i) => ({ sira: i }));
}

// 1. Donus degerini yok sayan uretici
const a = new YavasYazici();
for (const kayit of kayitlar(50)) a.write(kayit);
console.log('yok sayan  -> tamponda bekleyen kayit:', a.writableLength);
a.end();

// 2. Donus degerini dinleyen uretici
const b = new YavasYazici();
let enYuksek = 0;
for (const kayit of kayitlar(50)) {
  if (!b.write(kayit)) {
    enYuksek = Math.max(enYuksek, b.writableLength);
    await once(b, 'drain');
  }
}
console.log('dinleyen   -> tamponda gorulen en yuksek deger:', enYuksek);
b.end();
node geri-basinc.mjs
yok sayan  -> tamponda bekleyen kayit: 50
dinleyen   -> tamponda gorulen en yuksek deger: 4

İki sayı arasındaki fark, geri basıncın tamamıdır. Dönüş değerini yok sayan üretici elli kaydın tamamını belleğe yığdı; eşik 4 olmasına rağmen tampon 50’ye çıktı. Dönüş değerini dinleyen üretici tamponu hiç 4’ün üstüne çıkarmadı.

Bu farkın üretimdeki karşılığı bellek tüketimidir. Yavaş bir istemciye hızlı üretilen veriyi yazan bir sunucu, dönüş değerini yok sayarsa istemci başına sınırsız bellek biriktirir; birkaç yavaş istemci süreci bellek sınırına dayar.

İyi haber şudur: pipeline ve pipe bu denetimi kendiliğinden yapar; for await döngüsü de öyle. Dönüş değerini elle yönetmek yalnızca akışları doğrudan sürdüğünüz yerlerde gerekir.

Satır tabanlı okuma için ayrıca yerleşik bir kolaylık vardır: node:readline modülünün arayüzü, bir okunabilir akışı satır satır yinelenebilir hâle getirir ve yukarıdaki bölme mantığını üstlenir. Kendi dönüştürücünüzü yazmak, kaydın satırdan başka bir şeyle sınırlandığı biçimlerde gerekir.

Özet

  • Akışlar veriyi parça parça taşır; dört türü okunabilir, yazılabilir, çift yönlü ve dönüştürücüdür.
  • Parça sınırları kayıt sınırlarıyla örtüşmez; kaydı kurmak, yarım kalan veriyi saklayan bir dönüştürücünün işidir ve _flush kancası son kaydı kurtarır.
  • Nesne kipi, boru hattının bir noktasından sonra bayt yerine kayıt akmasını sağlar; dönüştürücünün iki tarafı ayrı yapılandırılabilir.
  • pipeline, zincirdeki bir hata durumunda bütün akışları yok eder ve hatayı tek yerden bildirir; pipe bu temizliği yapmaz.
  • write çağrısının false dönmesi geri basınç bildirimidir; yok sayıldığında tampon sınırsız büyür. pipeline, pipe ve for await bu denetimi üstlenir.

Sonraki Adım

Bu derste parçaların türü Buffer olarak göründü ve metne çevrilirken kodlama adı elle verildi. Parçanın bir karakterin ortasından bölünmesi durumunda ne olduğu ise sorulmadı. Sonraki ders ham baytların temsilini ele alır: arabellek nesnesinin ne olduğunu, çok baytlı karakterlerin parça sınırında nasıl bozulduğunu ve ikili bir kayıt biçiminin nasıl kodlanıp çözüldüğünü inceler.

İlerlemeni kaydetmek ve not almak için Giriş yap

Notlarım

Not almak için giriş yapmalısın.

Aramak için yazmaya başlayın.

↑↓ Esc gezin · aç · kapat