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
_flushkancası 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;pipebu temizliği yapmaz.writeçağrısınınfalsedönmesi geri basınç bildirimidir; yok sayıldığında tampon sınırsız büyür.pipeline,pipevefor awaitbu 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.