System.IO.Pipelines

System.IO.Pipelines adalah pustaka yang dirancang untuk membuat I/O berkinerja tinggi di .NET lebih mudah. Paket ini menargetkan .NET Standard untuk kompatibilitas luas, .NET Framework, dan .NET modern. Dalam versi .NET modern, System.IO.Pipelines disertakan dalam kerangka kerja bersama dan tidak memerlukan paket NuGet terpisah.

Pustaka juga tersedia sebagai paket System.IO.Pipelines NuGet.

Masalah apa yang dipecahkan oleh System.IO.Pipelines?

Aplikasi yang mengurai data streaming terdiri dari kode boilerplate yang memiliki banyak alur kode khusus dan tidak biasa. Kode template dan kode kasus khusus rumit dan sulit untuk dipertahankan.

System.IO.Pipelines dirancang untuk:

  • Memiliki penguraian data streaming dengan performa tinggi.
  • Mengurangi kompleksitas kode.

Kode ini khas untuk server TCP yang menerima pesan yang dibatasi baris (dibatasi oleh '\n') dari klien:

async Task ProcessLinesAsync(NetworkStream stream)
{
    var buffer = new byte[1024];
    await stream.ReadAsync(buffer, 0, buffer.Length);

    // Process a single line from the buffer
    ProcessLine(buffer);
}

Kode sebelumnya memiliki beberapa masalah:

  • Seluruh pesan (akhir baris) mungkin tidak diterima dalam satu panggilan ke ReadAsync.
  • Kode ini mengabaikan hasil stream.ReadAsync. stream.ReadAsync mengembalikan berapa banyak data yang dibaca.
  • Ini tidak menangani kasus di mana beberapa baris dibaca dalam satu panggilan ReadAsync.
  • Memalokasikan array byte dengan setiap pembacaan.

Untuk memperbaiki masalah sebelumnya, buat perubahan ini:

  • Buffer data masuk hingga baris baru ditemukan.

  • Analisis semua baris yang dikembalikan dalam buffer.

  • Garis mungkin lebih besar dari 1 KB (1024 byte). Kode perlu mengubah ukuran buffer input hingga pembatas ditemukan agar sesuai dengan baris lengkap di dalam buffer.

    • Jika buffer diubah ukurannya, lebih banyak salinan buffer akan dibuat saat muncul garis-garis yang lebih panjang di input.
    • Untuk mengurangi ruang yang terbuang, padatkan buffer yang digunakan untuk membaca baris.
  • Mempertimbangkan untuk menggunakan pengumpulan buffer untuk menghindari alokasi memori berulang kali.

  • Kode ini membahas beberapa masalah ini:

async Task ProcessLinesAsync(NetworkStream stream)
{
    byte[] buffer = ArrayPool<byte>.Shared.Rent(1024);
    var bytesBuffered = 0;
    var bytesConsumed = 0;

    while (true)
    {
        // Calculate the amount of bytes remaining in the buffer.
        var bytesRemaining = buffer.Length - bytesBuffered;

        if (bytesRemaining == 0)
        {
            // Double the buffer size and copy the previously buffered data into the new buffer.
            var newBuffer = ArrayPool<byte>.Shared.Rent(buffer.Length * 2);
            Buffer.BlockCopy(buffer, 0, newBuffer, 0, buffer.Length);
            // Return the old buffer to the pool.
            ArrayPool<byte>.Shared.Return(buffer);
            buffer = newBuffer;
            bytesRemaining = buffer.Length - bytesBuffered;
        }

        var bytesRead = await stream.ReadAsync(buffer, bytesBuffered, bytesRemaining);
        if (bytesRead == 0)
        {
            // EOF
            break;
        }

        // Keep track of the amount of buffered bytes.
        bytesBuffered += bytesRead;
        var linePosition = -1;

        do
        {
            // Look for a EOL in the buffered data.
            linePosition = Array.IndexOf(buffer, (byte)'\n', bytesConsumed,
                                         bytesBuffered - bytesConsumed);

            if (linePosition >= 0)
            {
                // Calculate the length of the line based on the offset.
                var lineLength = linePosition - bytesConsumed;

                // Process the line.
                ProcessLine(buffer, bytesConsumed, lineLength);

                // Move the bytesConsumed to skip past the line consumed (including \n).
                bytesConsumed += lineLength + 1;
            }
        }
        while (linePosition >= 0);
    }
}

Kode sebelumnya rumit dan tidak mengatasi semua masalah yang diidentifikasi. Jaringan berkinerja tinggi biasanya berarti menulis kode kompleks untuk memaksimalkan performa. System.IO.Pipelines dirancang untuk membuat penulisan jenis kode ini lebih mudah.

Pipa

Gunakan kelas Pipe untuk membuat pasangan PipeWriter/PipeReader. Semua data yang ditulis ke PipeWriter tersedia di dalam PipeReader:

var pipe = new Pipe();
PipeReader reader = pipe.Reader;
PipeWriter writer = pipe.Writer;

Penggunaan dasar pipa

async Task ProcessLinesAsync(Socket socket)
{
    var pipe = new Pipe();
    Task writing = FillPipeAsync(socket, pipe.Writer);
    Task reading = ReadPipeAsync(pipe.Reader);

    await Task.WhenAll(reading, writing);
}

async Task FillPipeAsync(Socket socket, PipeWriter writer)
{
    const int minimumBufferSize = 512;

    while (true)
    {
        // Allocate at least 512 bytes from the PipeWriter.
        Memory<byte> memory = writer.GetMemory(minimumBufferSize);
        try
        {
            int bytesRead = await socket.ReceiveAsync(memory, SocketFlags.None);
            if (bytesRead == 0)
            {
                break;
            }
            // Tell the PipeWriter how much was read from the Socket.
            writer.Advance(bytesRead);
        }
        catch (Exception ex)
        {
            LogError(ex);
            break;
        }

        // Make the data available to the PipeReader.
        FlushResult result = await writer.FlushAsync();

        if (result.IsCompleted)
        {
            break;
        }
    }

     // By completing PipeWriter, tell the PipeReader that there's no more data coming.
    await writer.CompleteAsync();
}

async Task ReadPipeAsync(PipeReader reader)
{
    while (true)
    {
        ReadResult result = await reader.ReadAsync();
        ReadOnlySequence<byte> buffer = result.Buffer;

        while (TryReadLine(ref buffer, out ReadOnlySequence<byte> line))
        {
            // Process the line.
            ProcessLine(line);
        }

        // Tell the PipeReader how much of the buffer has been consumed.
        reader.AdvanceTo(buffer.Start, buffer.End);

        // Stop reading if there's no more data coming.
        if (result.IsCompleted)
        {
            break;
        }
    }

    // Mark the PipeReader as complete.
    await reader.CompleteAsync();
}

bool TryReadLine(ref ReadOnlySequence<byte> buffer, out ReadOnlySequence<byte> line)
{
    // Look for a EOL in the buffer.
    SequencePosition? position = buffer.PositionOf((byte)'\n');

    if (position == null)
    {
        line = default;
        return false;
    }

    // Skip the line + the \n.
    line = buffer.Slice(0, position.Value);
    buffer = buffer.Slice(buffer.GetPosition(1, position.Value));
    return true;
}

Dua perulangan menangani pembacaan dan penulisan:

  • FillPipeAsync membaca dari Socket dan menulis ke PipeWriter.
  • ReadPipeAsync membaca dari PipeReader dan mengurai baris masuk.

Tidak ada buffer eksplisit yang dialokasikan. Semua manajemen buffer didelegasikan ke implementasi PipeReader dan PipeWriter. Mendelegasikan manajemen buffer memudahkan penggunaan kode untuk hanya berfokus pada logika bisnis.

Di perulangan pertama:

Dalam perulangan kedua, PipeReader mengonsumsi buffer yang ditulis oleh PipeWriter. Buffer berasal dari soket. Panggilan ke PipeReader.ReadAsync:

  • Mengembalikan ReadResult yang berisi dua informasi penting:

    • Data yang dibaca dalam bentuk ReadOnlySequence<T>.
    • Boolean IsCompleted yang menunjukkan apakah akhir data (EOF) telah tercapai.

Setelah menemukan pemisah akhir baris (EOL) dan mengurai baris:

  • Logika mengolah buffer untuk melewatkan bagian yang sudah diproses.
  • PipeReader.AdvanceTo dipanggil untuk memberi tahu PipeReader seberapa banyak data yang telah digunakan dan diperiksa.

Loop pembaca dan penulis diakhiri dengan memanggil PipeReader.Complete dan PipeWriter.Complete. Panggilan Complete melepaskan memori yang terkait dengan Pipe yang dialokasikan.

Tekanan balik dan pengendalian aliran

Idealnya, membaca dan mengurai bekerja sama:

  • Utas baca mengonsumsi data dari jaringan dan memasukkannya ke dalam buffer.
  • Thread penguraian bertanggung jawab untuk membangun struktur data yang sesuai.

Biasanya, penguraian membutuhkan lebih banyak waktu daripada hanya menyalin blok data dari jaringan:

  • Thread pembacaan lebih cepat dari thread penguraian.
  • Utas baca harus memperlambat kecepatan atau mengalokasikan lebih banyak memori untuk menyimpan data dalam proses penguraian.

Untuk performa optimal, terdapat keseimbangan antara jeda yang sering dan mengalokasikan lebih banyak memori.

Untuk mengatasi masalah sebelumnya, Pipe memiliki dua pengaturan untuk mengontrol aliran data:

Diagram yang menampilkan ResumeWriterThreshold dan PauseWriterThreshold

PipeWriter.FlushAsync:

  • Mengembalikan ValueTask<FlushResult> yang tidak lengkap saat jumlah data dalam Pipe melampaui PauseWriterThreshold.
  • Selesaikan ValueTask<FlushResult> ketika menjadi lebih rendah dari ResumeWriterThreshold.

Dua nilai digunakan untuk mencegah pengulangan cepat, yang dapat terjadi jika hanya satu nilai digunakan.

Contoh

// The Pipe will start returning incomplete tasks from FlushAsync until
// the reader examines at least 5 bytes.
var options = new PipeOptions(pauseWriterThreshold: 10, resumeWriterThreshold: 5);
var pipe = new Pipe(options);

PipeScheduler

Biasanya saat menggunakan async dan await, kode asinkron dilanjutkan pada TaskScheduler atau saat ini SynchronizationContext.

Saat melakukan I/O, penting untuk memiliki kontrol halus atas tempat I/O dilakukan. Kontrol ini memungkinkan memanfaatkan cache CPU secara efektif. Cache yang efisien sangat penting untuk aplikasi berkinerja tinggi seperti server web. PipeScheduler menyediakan kontrol di mana panggilan balik asinkron berjalan. Secara default:

  • Saat ini SynchronizationContext digunakan.
  • Jika tidak ada SynchronizationContext, maka kumpulan utas digunakan untuk menjalankan panggilan balik.
public static void Main(string[] args)
{
    var writeScheduler = new SingleThreadPipeScheduler();
    var readScheduler = new SingleThreadPipeScheduler();

    // Tell the Pipe what schedulers to use and disable the SynchronizationContext.
    var options = new PipeOptions(readerScheduler: readScheduler,
                                  writerScheduler: writeScheduler,
                                  useSynchronizationContext: false);
    var pipe = new Pipe(options);
}

// This is a sample scheduler that async callbacks on a single dedicated thread.
public class SingleThreadPipeScheduler : PipeScheduler
{
    private readonly BlockingCollection<(Action<object> Action, object State)> _queue =
     new BlockingCollection<(Action<object> Action, object State)>();
    private readonly Thread _thread;

    public SingleThreadPipeScheduler()
    {
        _thread = new Thread(DoWork);
        _thread.Start();
    }

    private void DoWork()
    {
        foreach (var item in _queue.GetConsumingEnumerable())
        {
            item.Action(item.State);
        }
    }

    public override void Schedule(Action<object?> action, object? state)
    {
        if (state is not null)
        {
            _queue.Add((action, state));
        }
        // else log the fact that _queue.Add was not called.
    }
}

PipeScheduler.ThreadPool adalah implementasi PipeScheduler yang mengantrekan panggilan balik ke kumpulan utas. PipeScheduler.ThreadPool adalah default dan umumnya pilihan terbaik. PipeScheduler.Inline dapat menyebabkan konsekuensi yang tidak diinginkan seperti kebuntuan.

Reset pipa

Menggunakan kembali objek Pipe sering kali efisien. Untuk mengatur ulang pipa, panggil PipeReaderReset saat PipeReader dan PipeWriter selesai.

PipeReader

PipeReader mengelola memori atas nama pemanggil. Selalu panggil PipeReader.AdvanceTo setelah memanggil PipeReader.ReadAsync. Ini memberi tahu PipeReader kapan penelepon selesai menggunakan memori sehingga dapat dilacak. Yang ReadOnlySequence<T> dikembalikan dari PipeReader.ReadAsync hanya valid hingga panggilan ke PipeReader.AdvanceTo. Hal ini ilegal untuk digunakan ReadOnlySequence<T> setelah memanggil PipeReader.AdvanceTo.

PipeReader.AdvanceTo memiliki dua argumen SequencePosition:

  • Argumen pertama menentukan berapa banyak memori yang digunakan.
  • Argumen kedua menentukan berapa banyak buffer yang diamati.

Menandai data sebagai telah digunakan berarti bahwa pipa dapat mengembalikan memori ke kumpulan buffer dasar. Menandai data sebagai diamati mengontrol apa yang akan dilakukan panggilan PipeReader.ReadAsync berikutnya. Menandai semuanya seperti yang diamati berarti bahwa panggilan berikutnya ke PipeReader.ReadAsync tidak akan kembali sampai ada lebih banyak data yang ditulis ke pipa. Nilai lain membuat panggilan berikutnya ke PipeReader.ReadAsync segera kembali dengan data yang diamati dan yang belum diamati, tetapi bukan data yang telah dikonsumsi.

Memahami skenario aliran data

Beberapa pola umum muncul saat membaca data streaming:

  • Mengingat aliran data, uraikan satu pesan.
  • Diberikan aliran data, memproses semua pesan yang tersedia.

Contoh-contoh ini menggunakan TryParseLines metode untuk mengurai pesan dari ReadOnlySequence<T>. TryParseLines mengurai satu pesan dan memperbarui buffer input untuk menghapus pesan yang telah diurai dari buffer. TryParseLines bukan bagian dari .NET; ini adalah metode yang ditulis pengguna yang digunakan di bagian berikut.

bool TryParseLines(ref ReadOnlySequence<byte> buffer, out Message message);

Membaca satu pesan

Kode ini membaca satu pesan dari PipeReader dan mengembalikannya ke pemanggil.

async ValueTask<Message?> ReadSingleMessageAsync(PipeReader reader,
 CancellationToken cancellationToken = default)
{
    while (true)
    {
        ReadResult result = await reader.ReadAsync(cancellationToken);
        ReadOnlySequence<byte> buffer = result.Buffer;

        // In the event that no message is parsed successfully, mark consumed
        // as nothing and examined as the entire buffer.
        SequencePosition consumed = buffer.Start;
        SequencePosition examined = buffer.End;

        try
        {
            if (TryParseLines(ref buffer, out Message message))
            {
                // A single message was successfully parsed so mark the start of the
                // parsed buffer as consumed. TryParseLines trims the buffer to
                // point to the data after the message was parsed.
                consumed = buffer.Start;

                // Examined is marked the same as consumed here, so the next call
                // to ReadSingleMessageAsync will process the next message if there's
                // one.
                examined = consumed;

                return message;
            }

            // There's no more data to be processed.
            if (result.IsCompleted)
            {
                if (buffer.Length > 0)
                {
                    // The message is incomplete and there's no more data to process.
                    throw new InvalidDataException("Incomplete message.");
                }

                break;
            }
        }
        finally
        {
            reader.AdvanceTo(consumed, examined);
        }
    }

    return null;
}

Kode sebelumnya:

  • Mengurai satu pesan.
  • Memperbarui yang dikonsumsi SequencePosition dan diperiksa SequencePosition untuk menunjuk ke awal buffer input yang dipangkas.

Dua argumen SequencePosition diperbarui karena TryParseLines menghapus pesan yang diurai dari buffer input. Umumnya, saat mengurai satu pesan dari buffer, posisi yang diperiksa harus menjadi salah satu hal berikut:

  • Akhir pesan.
  • Akhir dari buffer yang diterima jika tidak ada pesan yang ditemukan.

Kasus pesan tunggal memiliki potensi kesalahan paling besar. Meneruskan nilai yang salah ke diperiksa dapat mengakibatkan pengecualian kehabisan memori atau perulangan tak terbatas. Untuk informasi selengkapnya, lihat bagian Masalah umum PipeReader di artikel ini.

Penting

ReadSingleMessageAsync tidak memanggil PipeReader.CompleteAsync. Pemanggil bertanggung jawab untuk menyelesaikan PipeReader. Memanggil PipeReader.CompleteAsync dalam ReadSingleMessageAsync menandakan bahwa tidak ada data lebih lanjut yang dapat dibaca, yang mencegah pembacaan pesan berikutnya.

Membaca beberapa pesan

Kode ini membaca semua pesan dari PipeReader dan memanggil ProcessMessageAsync pada masing-masing pesan.

async Task ProcessMessagesAsync(PipeReader reader, CancellationToken cancellationToken = default)
{
    try
    {
        while (true)
        {
            ReadResult result = await reader.ReadAsync(cancellationToken);
            ReadOnlySequence<byte> buffer = result.Buffer;

            try
            {
                // Process all messages from the buffer, modifying the input buffer on each
                // iteration.
                while (TryParseLines(ref buffer, out Message message))
                {
                    await ProcessMessageAsync(message);
                }

                // There's no more data to be processed.
                if (result.IsCompleted)
                {
                    if (buffer.Length > 0)
                    {
                        // The message is incomplete and there's no more data to process.
                        throw new InvalidDataException("Incomplete message.");
                    }
                    break;
                }
            }
            finally
            {
                // Since all messages in the buffer are being processed, you can use the
                // remaining buffer's Start and End position to determine consumed and examined.
                reader.AdvanceTo(buffer.Start, buffer.End);
            }
        }
    }
    finally
    {
        await reader.CompleteAsync();
    }
}

Karena ProcessMessagesAsync memiliki perulangan pembacaan pesan lengkap, itu memanggil PipeReader.CompleteAsync ketika selesai. Tidak seperti kasus pesan tunggal, pemanggil tidak perlu menyelesaikan proses pembacaan. ProcessMessagesAsync mengambil kepemilikan penuh atas masa berlaku PipeReader.

Pembatalan

PipeReader.ReadAsync:

  • Mendukung melewati CancellationToken.
  • Melempar OperationCanceledException jika CancellationToken dibatalkan ketika proses membaca sedang berlangsung.
  • Mendukung cara untuk membatalkan operasi baca saat ini melalui PipeReader.CancelPendingRead, yang menghindari kemunculan pengecualian. Panggilan PipeReader.CancelPendingRead akan mengembalikan PipeReader.ReadAsync dengan ReadResult yang diatur ke IsCanceled pada panggilan true saat ini atau berikutnya. Ini berguna untuk menghentikan perulangan baca yang ada dengan cara yang tidak merusak dan tidak luar biasa.
private PipeReader reader;

public MyConnection(PipeReader reader)
{
    this.reader = reader;
}

public void Abort()
{
    // Cancel the pending read so the process loop ends without an exception.
    reader.CancelPendingRead();
}

public async Task ProcessMessagesAsync()
{
    try
    {
        while (true)
        {
            ReadResult result = await reader.ReadAsync();
            ReadOnlySequence<byte> buffer = result.Buffer;

            try
            {
                if (result.IsCanceled)
                {
                    // The read was canceled. You can quit without reading the existing data.
                    break;
                }

                // Process all messages from the buffer, modifying the input buffer on each
                // iteration.
                while (TryParseLines(ref buffer, out Message message))
                {
                    await ProcessMessageAsync(message);
                }

                // There's no more data to be processed.
                if (result.IsCompleted)
                {
                    break;
                }
            }
            finally
            {
                // Since all messages in the buffer are being processed, you can use the
                // remaining buffer's Start and End position to determine consumed and examined.
                reader.AdvanceTo(buffer.Start, buffer.End);
            }
        }
    }
    finally
    {
        await reader.CompleteAsync();
    }
}

Masalah umum PipeReader

  • Meneruskan nilai yang salah ke consumed atau examined dapat mengakibatkan pembacaan ulang data yang sudah terbaca.

  • Menggunakan buffer.End sebagai yang diperiksa dapat mengakibatkan:

    • Data terhenti
    • Pengecualian out-of-memory (OOM) yang pada akhirnya terjadi jika data tidak dikonsumsi. Misalnya, PipeReader.AdvanceTo(position, buffer.End) saat memproses satu pesan pada satu waktu dari buffer.
  • Meneruskan nilai yang salah ke consumed atau examined mungkin mengakibatkan perulangan tak terbatas. Misalnya, PipeReader.AdvanceTo(buffer.Start) jika buffer.Start belum berubah menyebabkan panggilan PipeReader.ReadAsync berikutnya segera kembali sebelum data baru tiba.

  • Meneruskan nilai yang salah ke consumed atau examined mungkin mengakibatkan buffering tak terbatas (akhirnya OOM).

  • Menggunakan ReadOnlySequence<T> setelah memanggil PipeReader.AdvanceTo dapat mengakibatkan korupsi memori (penggunaan setelah pembebasan).

  • Gagal memanggil Complete/CompleteAsync dapat mengakibatkan kebocoran memori.

  • Memeriksa ReadResult.IsCompleted dan keluar dari logika pembacaan sebelum memproses buffer mengakibatkan hilangnya data. Kondisi penghentian perulangan harus didasarkan pada ReadResult.Buffer.IsEmpty dan ReadResult.IsCompleted. Melakukan ini dengan tidak benar dapat mengakibatkan perulangan tak terbatas.

Kode bermasalah

Kehilangan data

ReadResult dapat mengembalikan segmen akhir data saat IsCompleted diatur ke true. Tidak membaca data tersebut sebelum keluar dari perulangan baca menghasilkan kehilangan data.

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

Environment.FailFast("This code is terrible, don't use it!");
while (true)
{
    ReadResult result = await reader.ReadAsync(cancellationToken);
    ReadOnlySequence<byte> dataLossBuffer = result.Buffer;

    if (result.IsCompleted)
        break;

    Process(ref dataLossBuffer, out Message message);

    reader.AdvanceTo(dataLossBuffer.Start, dataLossBuffer.End);
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

Loop tak terbatas

Logika berikut dapat mengakibatkan perulangan tak terbatas jika Result.IsCompleted ada true tetapi tidak pernah ada pesan lengkap dalam buffer.

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

Environment.FailFast("This code is terrible, don't use it!");
while (true)
{
    ReadResult result = await reader.ReadAsync(cancellationToken);
    ReadOnlySequence<byte> infiniteLoopBuffer = result.Buffer;
    if (result.IsCompleted && infiniteLoopBuffer.IsEmpty)
        break;

    Process(ref infiniteLoopBuffer, out Message message);

    reader.AdvanceTo(infiniteLoopBuffer.Start, infiniteLoopBuffer.End);
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

Berikut adalah bagian lain dari kode dengan masalah yang sama. Hal ini memeriksa buffer yang tidak kosong sebelum memeriksa ReadResult.IsCompleted. Karena dalam else if, itu akan terus berulang selamanya jika tidak ada pesan lengkap di buffer.

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

Environment.FailFast("This code is terrible, don't use it!");
while (true)
{
    ReadResult result = await reader.ReadAsync(cancellationToken);
    ReadOnlySequence<byte> infiniteLoopBuffer = result.Buffer;

    if (!infiniteLoopBuffer.IsEmpty)
        Process(ref infiniteLoopBuffer, out Message message);

    else if (result.IsCompleted)
        break;

    reader.AdvanceTo(infiniteLoopBuffer.Start, infiniteLoopBuffer.End);
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

Aplikasi yang tidak responsif

Panggilan PipeReader.AdvanceTo tanpa syarat dengan buffer.End dalam examined posisi dapat mengakibatkan aplikasi menjadi tidak responsif saat mengurai satu pesan. Panggilan berikutnya untuk PipeReader.AdvanceTo tidak akan kembali hingga:

  • Terdapat lebih banyak data yang ditulis ke pipa.
  • Dan data baru sebelumnya tidak diperiksa.

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

Environment.FailFast("This code is terrible, don't use it!");
while (true)
{
    ReadResult result = await reader.ReadAsync(cancellationToken);
    ReadOnlySequence<byte> hangBuffer = result.Buffer;

    Process(ref hangBuffer, out Message message);

    if (result.IsCompleted)
        break;

    reader.AdvanceTo(hangBuffer.Start, hangBuffer.End);

    if (message != null)
        return message;
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

Kehabisan Memori (OOM)

Dengan kondisi yang berikut, kode ini mempertahankan buffering hingga OutOfMemoryException terjadi.

  • Tidak ada ukuran pesan maksimum.
  • Data yang dikembalikan dari PipeReader tidak membuat pesan lengkap. Misalnya, sisi lain menulis pesan besar (misalnya, pesan 4 GB).

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

Environment.FailFast("This code is terrible, don't use it!");
while (true)
{
    ReadResult result = await reader.ReadAsync(cancellationToken);
    ReadOnlySequence<byte> thisCouldOutOfMemory = result.Buffer;

    Process(ref thisCouldOutOfMemory, out Message message);

    if (result.IsCompleted)
        break;

    reader.AdvanceTo(thisCouldOutOfMemory.Start, thisCouldOutOfMemory.End);

    if (message != null)
        return message;
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

Kerusakan Memori

Saat menulis fungsi pembantu yang membaca buffer, salin payload yang dikembalikan sebelum memanggil Advance. Contoh berikut mengembalikan memori yang Pipe telah dibuang dan mungkin menggunakannya kembali untuk operasi berikutnya (baca/tulis).

Peringatan

JANGAN gunakan kode berikut. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel berikut disediakan untuk menjelaskan masalah umum PipeReader.

public class Message
{
    public ReadOnlySequence<byte> CorruptedPayload { get; set; }
}
    Environment.FailFast("This code is terrible, don't use it!");
    Message message = null;

    while (true)
    {
        ReadResult result = await reader.ReadAsync(cancellationToken);
        ReadOnlySequence<byte> buffer = result.Buffer;

        ReadHeader(ref buffer, out int length);

        if (length <= buffer.Length)
        {
            message = new Message
            {
                // Slice the payload from the existing buffer
                CorruptedPayload = buffer.Slice(0, length)
            };

            buffer = buffer.Slice(length);
        }

        if (result.IsCompleted)
            break;

        reader.AdvanceTo(buffer.Start, buffer.End);

        if (message != null)
        {
            // This code is broken since reader.AdvanceTo() was called with a position *after* the buffer
            // was captured.
            break;
        }
    }

    return message;
}

Peringatan

JANGAN gunakan kode sebelumnya. Menggunakan sampel ini akan mengakibatkan kehilangan data, macet, masalah keamanan, dan TIDAK boleh disalin. Sampel sebelumnya disediakan untuk menjelaskan masalah umum PipeReader.

PipeWriter

PipeWriter mengelola buffer untuk menulis atas nama pemanggil. PipeWriter mengimplementasikan IBufferWriter<byte>. IBufferWriter<byte> menyediakan akses ke buffer untuk melakukan penulisan tanpa salinan buffer tambahan.

async Task WriteHelloAsync(PipeWriter writer, CancellationToken cancellationToken = default)
{
    // Request at least 5 bytes from the PipeWriter.
    Memory<byte> memory = writer.GetMemory(5);

    // Write directly into the buffer.
    int written = Encoding.ASCII.GetBytes("Hello".AsSpan(), memory.Span);

    // Tell the writer how many bytes were written.
    writer.Advance(written);

    await writer.FlushAsync(cancellationToken);
}

Kode sebelumnya:

  • Meminta buffer setidaknya 5 byte dari PipeWriter menggunakan GetMemory.
  • Menulis byte untuk string ASCII "Hello" ke Memory<byte> yang dikembalikan.
  • Panggilan Advance untuk menunjukkan berapa banyak byte yang ditulis ke buffer.
  • Menghapus PipeWriter, yang mengirim byte ke perangkat yang mendasar.

Metode penulisan sebelumnya menggunakan buffer yang disediakan oleh PipeWriter. Selain itu, dapat menggunakan PipeWriter.WriteAsync, yang:

async Task WriteHelloAsync(PipeWriter writer, CancellationToken cancellationToken = default)
{
    byte[] helloBytes = Encoding.ASCII.GetBytes("Hello");

    // Write helloBytes to the writer, there's no need to call Advance here
    // (Write does that).
    await writer.WriteAsync(helloBytes, cancellationToken);
}

Pembatalan

FlushAsync mendukung melewati CancellationToken. Meneruskan CancellationToken menyebabkan OperationCanceledException apabila token dibatalkan saat flush masih tertunda. PipeWriter.FlushAsync mendukung pembatalan operasi flush saat ini melalui PipeWriter.CancelPendingFlush tanpa memunculkan pengecualian. Panggilan PipeWriter.CancelPendingFlush menyebabkan panggilan saat ini atau berikutnya ke PipeWriter.FlushAsync atau PipeWriter.WriteAsync akan mengembalikan FlushResult dengan IsCanceled yang diatur ke true. Ini berguna untuk menghentikan flush yang menghasilkan dengan cara yang tidak merusak dan tidak luar biasa.

Masalah umum PipeWriter

  • GetSpan dan GetMemory mengembalikan buffer dengan setidaknya jumlah memori yang diminta. Jangan asumsikan ukuran buffer yang tepat.
  • Panggilan berturut-turut tidak dijamin menghasilkan buffer yang sama atau ukuran buffer yang sama.
  • Buffer baru harus diminta setelah memanggil Advance untuk terus menulis lebih banyak data. Buffer yang sudah diperoleh sebelumnya tidak bisa ditulisi.
  • Panggilan GetMemory atau GetSpan saat ada panggilan yang tidak lengkap ke FlushAsync tidak aman.
  • Memanggil Complete atau CompleteAsync saat ada data yang belum dihapus dapat mengakibatkan kerusakan memori.

Tips untuk PipeReader dan PipeWriter

Gunakan tips ini untuk berhasil menggunakan System.IO.Pipelines kelas:

  • Selalu selesaikan PipeReader dan PipeWriter, termasuk pengecualian jika berlaku.
  • Selalu panggil PipeReader.AdvanceTo setelah memanggil PipeReader.ReadAsync.
  • Secara berkala awaitPipeWriter.FlushAsync saat menulis, dan selalu periksa FlushResult.IsCompleted. Batalkan penulisan jika IsCompleted adalah true, seperti yang menunjukkan pembaca selesai dan tidak lagi peduli tentang apa yang ditulis.
  • Panggil PipeWriter.FlushAsync setelah menulis sesuatu yang Anda inginkan bisa diakses oleh PipeReader.
  • Jangan panggil FlushAsync jika pembaca tidak dapat memulai sampai FlushAsync selesai, karena itu dapat menyebabkan kebuntuan.
  • Pastikan bahwa hanya satu konteks yang "memiliki" PipeReader atau PipeWriter atau mengaksesnya. Jenis ini tidak aman untuk utas.
  • Jangan pernah mengakses ReadResult.Buffer setelah memanggil PipeReader.AdvanceTo atau menyelesaikan PipeReader.

IDuplexPipe

IDuplexPipe adalah kontrak bagi tipe yang mendukung pembacaan dan penulisan. Misalnya, koneksi jaringan diwakili oleh IDuplexPipe.

Tidak seperti Pipe, yang berisi PipeReader dan PipeWriter, IDuplexPipe mewakili satu sisi koneksi dupleks penuh. Apa yang Anda tulis ke PipeWriter tidak akan dibaca dari PipeReader.

Aliran

Saat membaca atau menulis data aliran, Anda biasanya membaca data menggunakan de-serializer dan menulis data menggunakan serializer. Sebagian besar API aliran baca dan tulis ini memiliki parameter Stream. Untuk mempermudah integrasi dengan API yang ada ini, PipeReader dan PipeWriter mengekspos metode AsStream. AsStream mengembalikan implementasi Stream di sekitar PipeReader atau PipeWriter.

Contoh aliran

Buat instance PipeReader dan PipeWriter menggunakan metode statis Create yang diberikan sebuah objek Stream dan opsi pembuatan terkait yang bersifat opsional.

StreamPipeReaderOptions izinkan kontrol atas pembuatan instans PipeReader dengan parameter berikut:

StreamPipeWriterOptions izinkan kontrol atas pembuatan instans PipeWriter dengan parameter berikut:

Penting

Saat membuat PipeReader dan PipeWriter instans menggunakan Create metode , pertimbangkan Stream masa pakai objek. Jika Anda memerlukan akses ke aliran setelah pembaca atau penulis selesai dengannya, atur LeaveOpen flag ke true dalam opsi pembuatan. Jika tidak, aliran ditutup.

Kode ini menunjukkan cara membuat instance PipeReader dan PipeWriter menggunakan metode Create dari aliran.

using System.Buffers;
using System.IO.Pipelines;
using System.Text;

class Program
{
    static async Task Main()
    {
        using var stream = File.OpenRead("lorem-ipsum.txt");

        var reader = PipeReader.Create(stream);
        var writer = PipeWriter.Create(
            Console.OpenStandardOutput(), 
            new StreamPipeWriterOptions(leaveOpen: true));

        WriteUserCancellationPrompt();

        var processMessagesTask = ProcessMessagesAsync(reader, writer);
        var userCanceled = false;
        var cancelProcessingTask = Task.Run(() =>
        {
            while (char.ToUpperInvariant(Console.ReadKey().KeyChar) != 'C')
            {
                WriteUserCancellationPrompt();
            }

            userCanceled = true;

            // No exceptions thrown
            reader.CancelPendingRead();
            writer.CancelPendingFlush();
        });

        await Task.WhenAny(cancelProcessingTask, processMessagesTask);

        Console.WriteLine(
            $"\n\nProcessing {(userCanceled ? "cancelled" : "completed")}.\n");
    }

    static void WriteUserCancellationPrompt() =>
        Console.WriteLine("Press 'C' to cancel processing...\n");

    static async Task ProcessMessagesAsync(
        PipeReader reader,
        PipeWriter writer)
    {
        try
        {
            while (true)
            {
                ReadResult readResult = await reader.ReadAsync();
                ReadOnlySequence<byte> buffer = readResult.Buffer;

                try
                {
                    if (readResult.IsCanceled)
                    {
                        break;
                    }

                    if (TryParseLines(ref buffer, out string message))
                    {
                        FlushResult flushResult =
                            await WriteMessagesAsync(writer, message);

                        if (flushResult.IsCanceled || flushResult.IsCompleted)
                        {
                            break;
                        }
                    }

                    if (readResult.IsCompleted)
                    {
                        if (!buffer.IsEmpty)
                        {
                            throw new InvalidDataException("Incomplete message.");
                        }
                        break;
                    }
                }
                finally
                {
                    reader.AdvanceTo(buffer.Start, buffer.End);
                }
            }
        }
        catch (Exception ex)
        {
            Console.Error.WriteLine(ex);
        }
        finally
        {
            await reader.CompleteAsync();
            await writer.CompleteAsync();
        }
    }

    static bool TryParseLines(
        ref ReadOnlySequence<byte> buffer,
        out string message)
    {
        SequencePosition? position;
        StringBuilder outputMessage = new();

        while(true)
        {
            position = buffer.PositionOf((byte)'\n');

            if (!position.HasValue)
                break;

            outputMessage.Append(Encoding.ASCII.GetString(buffer.Slice(buffer.Start, position.Value)))
                        .AppendLine();

            buffer = buffer.Slice(buffer.GetPosition(1, position.Value));
        };

        message = outputMessage.ToString();
        return message.Length != 0;
    }

    static ValueTask<FlushResult> WriteMessagesAsync(
        PipeWriter writer,
        string message) =>
        writer.WriteAsync(Encoding.ASCII.GetBytes(message));
}

Aplikasi ini menggunakan StreamReader untuk membaca file lorem-ipsum.txt sebagai aliran, dan harus diakhir dengan baris kosong. FileStream diteruskan ke PipeReader.Create, yang membuat instans objek PipeReader. Aplikasi konsol kemudian meneruskan aliran output standarnya untuk PipeWriter.Create yang menggunakan Console.OpenStandardOutput(). Contoh mendukung pembatalan.