Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
System.IO.Pipelines je knihovna navržená tak, aby usnadnila vysoce výkonné vstupně-výstupní operace v .NET. Balíček cílí na .NET Standard pro zajištění široké kompatibility, rozhraní .NET Framework a moderního rozhraní .NET. V moderních verzích System.IO.Pipelines .NET je součástí sdílené architektury a nevyžaduje samostatný balíček NuGet.
Knihovna je také k dispozici jako balíček NuGet System.IO.Pipelines .
Jaký problém řeší System.IO.Pipelines?
Aplikace, které analyzují streamovaná data, se skládají z často používaného kódu, který má mnoho specializovaných a neobvyklých toků kódu. Šablonový kód a kód pro speciální případy je složitý a obtížně udržovatelný.
System.IO.Pipelines byl navržen tak, aby:
- Analýza streamovaných dat s vysokým výkonem
- Snižte složitost kódu.
Tento kód je typický pro server TCP, který přijímá řádkově oddělené zprávy (oddělené pomocí '\n') od klienta.
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);
}
Předchozí kód má několik problémů:
- Celá zpráva (konec řádku) nemusí být přijata v jednom volání
ReadAsync. - Ignoruje výsledek
stream.ReadAsync.stream.ReadAsyncvrátí, kolik dat bylo přečteno. - Nezpracuje případ, kdy se v jednom
ReadAsyncvolání čte více řádků. - Alokuje pole
bytepři každém čtení.
Pokud chcete vyřešit předchozí problémy, proveďte tyto změny:
Uložit příchozí data do vyrovnávací paměti, dokud není nalezen nový řádek.
Parsujte všechny řádky vrácené ve vyrovnávací paměti.
Čára může být větší než 1 kB (1024 bajtů). Kód musí změnit velikost vstupní vyrovnávací paměti, dokud nenajde oddělovač, aby se celý řádek vešel uvnitř vyrovnávací paměti.
- Pokud se velikost vyrovnávací paměti změní, vytvoří se více kopií vyrovnávací paměti, jakmile se ve vstupu zobrazí delší řádky.
- Pokud chcete omezit plýtvání místem, zahuštěte vyrovnávací paměť použitou pro čtení řádků.
Zvažte použití sdružování vyrovnávacích pamětí, abyste se vyhnuli opakovanému přidělování paměti.
Tento kód řeší některé z těchto problémů:
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);
}
}
Předchozí kód je složitý a neřeší všechny zjištěné problémy. Vysoce výkonné sítě obvykle znamenají psaní složitého kódu pro maximalizaci výkonu.
System.IO.Pipelines byl navržen tak, aby usnadnil psaní tohoto typu kódu.
Potrubí
Použijte třídu Pipe k vytvoření páru PipeWriter/PipeReader. Všechna data zapsaná do služby PipeWriter jsou k dispozici v :PipeReader
var pipe = new Pipe();
PipeReader reader = pipe.Reader;
PipeWriter writer = pipe.Writer;
Základní využití kanálu
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;
}
Dvě smyčky zpracovávají čtení a zápis:
-
FillPipeAsyncčte zSocketa zapisuje doPipeWriter. -
ReadPipeAsyncčte z příchozíchPipeReaderřádků a parsuje je.
Nejsou přiděleny žádné explicitní vyrovnávací paměti. Veškerá správa vyrovnávací paměti je delegována na implementace PipeReader a PipeWriter. Delegování správy vyrovnávacích pamětí usnadňuje využívání kódu výhradně na obchodní logiku.
V první smyčce:
- PipeWriter.GetMemory(Int32) je použita pro získání paměti z podkladového zapisovače.
-
PipeWriter.Advance(Int32) je volána, aby sdělila
PipeWriter, kolik dat bylo zapsáno do vyrovnávací paměti. -
PipeWriter.FlushAsync je volána pro zpřístupnění dat
PipeReader.
Ve druhé smyčce PipeReader využívá buffery zapsané pomocí PipeWriter. Vyrovnávací paměti pocházejí ze socketu. Toto volání:PipeReader.ReadAsync
ReadResult Vrátí hodnotu, která obsahuje dva důležité informace:
- Data, která byla načtena ve formě ReadOnlySequence<T>.
- Logická hodnota
IsCompleted, která označuje, jestli bylo dosaženo konce dat (EOF).
Po nalezení oddělovače konce řádku (EOL) a parsování řádku:
- Logika zpracovává vyrovnávací paměť tak, aby přeskočí, co už bylo zpracováno.
-
PipeReader.AdvanceTo je volána, aby informovala
PipeReader, kolik dat bylo spotřebováno a analyzováno.
Smyčky čtečky a zapisovače končí voláním PipeReader.Complete a PipeWriter.Complete. Volání Complete uvolní paměť přidělenou podkladovou složkou Pipe.
Zpětný tlak a řízení toku
V ideálním případě spolupracují čtení a analýza:
- Vlákno pro čtení spotřebovává data ze sítě a vkládá je do bufferů.
- Vlákno parsování zodpovídá za vytvoření odpovídajících datových struktur.
Analýza obvykle trvá déle než pouhé kopírování bloků dat ze sítě:
- Vlákno čtení předběhne vlákno zpracování.
- Čtecí vlákno musí buď zpomalit, nebo přidělit více paměti pro uložení dat pro vlákno zpracování.
Pro zajištění optimálního výkonu existuje rovnováhu mezi častými pozastaveními a přidělením více paměti.
Pokud chcete vyřešit předchozí problém, Pipe má dvě nastavení pro řízení toku dat:
- PauseWriterThreshold: Určuje, kolik dat by se mělo ukládat do vyrovnávací paměti před voláním na pozastavení FlushAsync.
- ResumeWriterThreshold: Určuje, kolik dat musí čtenář sledovat, než se obnoví volání na PipeWriter.FlushAsync.
- Vrátí neúplnou
ValueTask<FlushResult>, když množství dat vPipepřekročíPauseWriterThreshold. - Dokončí se, jakmile
ValueTask<FlushResult>bude nižší nežResumeWriterThreshold.
Dvě hodnoty se používají k prevenci rychlého cyklistiky, ke kterému může dojít, pokud se použije jedna hodnota.
Příklady
// 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);
Plánovač potrubí
Obvykle při použití async a await asynchronní kód pokračuje buď na TaskScheduler, nebo na aktuálním SynchronizationContext.
Při provádění vstupně-výstupních operací je důležité mít jemně odstupňovanou kontrolu nad tím, kde se vstupně-výstupní operace provádí. Tento ovládací prvek umožňuje efektivně využívat mezipaměti procesoru. Efektivní ukládání do mezipaměti je důležité pro vysoce výkonné aplikace, jako jsou webové servery. PipeScheduler poskytuje kontrolu nad tím, kde běží zpětná volání asynchronně. Standardně:
- Používá se současný SynchronizationContext.
- Pokud neexistuje
SynchronizationContext, používá fond vláken ke spouštění zpětných volání.
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 je PipeScheduler implementace, která zařazuje zpětná volání do vláknového fondu.
PipeScheduler.ThreadPool je výchozí a obecně nejlepší volbou.
PipeScheduler.Inline může způsobit nezamýšlené důsledky, jako jsou zablokování.
Resetování potrubí
Opakované využití objektu Pipe je často efektivní. Aby bylo možné obnovit potrubí, zavolejte PipeReaderReset po dokončení PipeReader i PipeWriter.
PipeReader
PipeReader spravuje paměť jménem volajícího.
Vždy zavolejte PipeReader.AdvanceTo po zavolání PipeReader.ReadAsync. Toto umožňuje PipeReader zjistit, kdy je volající hotový s pamětí, aby mohla být sledována. Hodnota ReadOnlySequence<T> vrácená z PipeReader.ReadAsync je platná pouze do volání PipeReader.AdvanceTo. Je nezákonné používat ReadOnlySequence<T> po volání PipeReader.AdvanceTo.
PipeReader.AdvanceTo přebírá dva SequencePosition argumenty:
- První argument určuje, kolik paměti bylo spotřebováno.
- Druhý argument určuje, kolik vyrovnávací paměti bylo zjištěno.
Označení dat jako spotřebovaných znamená, že potrubí může vrátit paměť zpět do vyrovnávací paměťové základny. Označení dat jako pozorovaných určuje, co provede další volání PipeReader.ReadAsync. Označení všeho jako dohlédnutého znamená, že další volání PipeReader.ReadAsync se nevrátí, dokud nebudou do potrubí zapsána další data. Jakákoli jiná hodnota způsobí, že se PipeReader.ReadAsync okamžitě vrátí s pozorovanými a nepozorovanými daty, ale ne s daty, která už byla spotřebována.
Scénáře čtení streamovaných dat
Při čtení streamovaných dat se objeví několik typických vzorů:
- Z datového proudu parsuje jednu zprávu.
- Pro analýzu datového proudu zpracujte všechny dostupné zprávy.
Tyto příklady používají metodu TryParseLines pro analýzu zpráv z objektu ReadOnlySequence<T>.
TryParseLines parsuje jednu zprávu a aktualizuje vstupní vyrovnávací paměť, aby ořízla analyzovanou zprávu z vyrovnávací paměti.
TryParseLines není součástí .NET; je to uživatelsky napsaná metoda používaná v následujících částech.
bool TryParseLines(ref ReadOnlySequence<byte> buffer, out Message message);
Čtení jedné zprávy
Tento kód přečte jednu zprávu z PipeReader, a poté ji vrátí volajícímu.
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;
}
Předchozí kód:
- Parsuje jednu zprávu.
- Aktualizuje spotřebovanou
SequencePositioni zkoumanouSequencePosition, aby ukazovaly na začátek oříznuté vstupní vyrovnávací paměti.
Tyto dva SequencePosition argumenty jsou aktualizovány, protože TryParseLines odebere analyzovanou zprávu ze vstupní vyrovnávací paměti. Obecně platí, že při analýze jedné zprávy z vyrovnávací paměti by zkoumaná pozice měla být jedna z následujících možností:
- Konec zprávy.
- Konec přijaté vyrovnávací paměti, pokud nebyla nalezena žádná zpráva.
Jeden případ zprávy má největší potenciál pro chyby. Předání nesprávných hodnot examined může vést k výjimce nedostatku paměti nebo nekonečné smyčce. Další informace naleznete v části Běžné problémy PipeReader v tomto článku.
Důležité
ReadSingleMessageAsync nevolá PipeReader.CompleteAsync. Volající je zodpovědný za dokončení hovoru PipeReader. Volání PipeReader.CompleteAsync uvnitř ReadSingleMessageAsync signalizuje, že již nelze číst žádná další data, což znemožňuje čtení následných zpráv.
Čtení více zpráv
Tento kód čte všechny zprávy z PipeReader a na každou volá ProcessMessageAsync.
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();
}
}
Protože ProcessMessagesAsync vlastní úplnou smyčku čtení zpráv, volá PipeReader.CompleteAsync po dokončení. Na rozdíl od případu s jednou zprávou nemusí volající dokončovat čtečku.
ProcessMessagesAsync převezme plnou odpovědnost za celou dobu životnosti PipeReader.
Zrušení
- Podporuje předání CancellationToken.
-
OperationCanceledException Vyvolá chybu, pokud
CancellationTokenje zrušena, zatímco čeká na čtení. - Podporuje způsob, jak zrušit aktuální operaci čtení prostřednictvím PipeReader.CancelPendingRead, což zabraňuje vyvolání výjimky. Volání
PipeReader.CancelPendingReadzpůsobí, že aktuální nebo příští volání PipeReader.ReadAsync vrátí ReadResult, který máIsCancelednastaveno natrue. To je užitečné pro zastavení stávající smyčky čtení nedestruktivním a nevýkonným způsobem.
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();
}
}
Běžné problémy s PipeReader
Předání nesprávných hodnot k
consumedaexaminedmůže vést ke čtení dat, která již byla přečtena.Předání
buffer.End, jak bylo zkoumáno, může mít za následek:- Pozastavená data
- Případná výjimka z nedostatku paměti (OOM), pokud data nejsou spotřebovávána. Například
PipeReader.AdvanceTo(position, buffer.End)při zpracování jedné zprávy najednou z vyrovnávací paměti.
Předání nesprávných hodnot do
consumedneboexaminedmůže vést k nekonečné smyčce. Pokud sePipeReader.AdvanceTo(buffer.Start)nezmění, způsobí, že příští volání PipeReader.ReadAsync se vrátí okamžitě, ještě než přijdou nová data.Předání nesprávných hodnot k
consumedneboexaminedmůže vést k nekonečnému ukládání do vyrovnávací paměti (případná chyba nedostatku paměti).Použití ReadOnlySequence<T> po volání PipeReader.AdvanceTo může vést k poškození paměti (užití po uvolnění).
Opomenutí volání Complete/CompleteAsync může vést k úniku paměti.
Kontrola ReadResult.IsCompleted a opuštění čtecí logiky před zpracováním pufru způsobí ztrátu dat. Podmínka ukončení smyčky by měla být založena na
ReadResult.Buffer.IsEmptyaReadResult.IsCompleted. Když to uděláte nesprávně, může to vést k nekonečné smyčce.
Problematický kód
❌ Ztráta dat
Může ReadResult vrátit poslední segment dat, pokud IsCompleted je nastavena na true. Nepřečtení dat před ukončením smyčky čtení vede ke ztrátě dat.
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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);
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
❌ Nekonečná smyčka
Následující logika může vést k nekonečné smyčce, pokud je Result.IsCompletedtrue, ale ve vyrovnávací paměti není nikdy úplná zpráva.
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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);
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
Tady je další část kódu se stejným problémem. Nejprve kontroluje, zda je buffer neprázdný, před kontrolou ReadResult.IsCompleted. Protože je v rámci else if, smyčka běží do nekonečna, pokud ve vyrovnávací paměti nikdy není úplná zpráva.
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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);
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
❌ Nereagující aplikace
Nepodmíněné volání PipeReader.AdvanceTo s buffer.End na pozici examined může vést k tomu, že při analýze jedné zprávy aplikace přestane reagovat. Další volání na PipeReader.AdvanceTo se nevrátí, dokud:
- Do kanálu se zapisuje další data.
- A nová data se ještě nezkoumala.
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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;
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
❌ Nedostatek paměti (OOM)
S následujícími podmínkami tento kód ukládá do vyrovnávací paměti, dokud nedojde k
- Neexistuje žádná maximální velikost zprávy.
- Data vrácená z
PipeReadernetvoří úplnou zprávu. Druhá strana například píše velkou zprávu (například zprávu o velikosti 4 GB).
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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;
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
❌ Poškození paměti
Při psaní pomocného programu, který čte z vyrovnávací paměti, zkopírujte libovolnou vrácenou datovou část před voláním Advance. Následující příklad vrátí paměť, kterou Pipe zahodila, a může ji znovu použít pro další operaci (čtení/zápis).
Upozornění
NEPOUŽÍVEJTE následující kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Následující ukázka je k dispozici k vysvětlení běžných problémů 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;
}
Upozornění
Nepoužívejte předchozí kód. Použití této ukázky způsobí ztrátu dat, zablokování, problémy se zabezpečením a nemělo by se kopírovat. Předchozí ukázka je k dispozici k vysvětlení běžných problémů PipeReader.
PipeWriter
PipeWriter spravuje vyrovnávací paměti pro zápis jménem volajícího.
PipeWriter implementuje IBufferWriter<byte>.
IBufferWriter<byte> poskytuje přístup k vyrovnávacím pamětím pro provádění zápisů bez dodatečných kopií vyrovnávací paměti.
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);
}
Předchozí kód:
- Z
PipeWriterpožaduje vyrovnávací paměť o velikosti alespoň 5 bajtů pomocí GetMemory. - Zapíše bajty pro řetězec
"Hello"ASCII do vrácenéMemory<byte>. - Volání Advance označuje, kolik bajtů bylo zapsáno do vyrovnávací paměti.
- Vyprázdní
PipeWriter, což odešle bajty do podkladového zařízení.
Metoda zápisu, která byla použita předtím, používá vyrovnávací paměti poskytované PipeWriter. Může také použít PipeWriter.WriteAsync, což:
- Zkopíruje existující vyrovnávací paměť do
PipeWriter. - Volání GetSpan, Advance podle potřeby a volání FlushAsync.
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);
}
Zrušení
FlushAsync podporuje předávání CancellationToken. Předání CancellationToken má za následek OperationCanceledException v případě, že je token zrušen, zatímco je vyprázdnění čekající.
PipeWriter.FlushAsync podporuje způsob, jak zrušit aktuální operaci vyprázdnění prostřednictvím PipeWriter.CancelPendingFlush bez vyvolání výjimky. Volání PipeWriter.CancelPendingFlush způsobí, že při aktuálním nebo příštím volání PipeWriter.FlushAsync nebo PipeWriter.WriteAsync se vrátí FlushResult, ve kterém je IsCanceled nastavena na true. To je užitečné pro zastavení vyprázdnění nedestruktivním a nevýkonným způsobem.
PipeWriter – běžné problémy
- GetSpan a GetMemory vracejí vyrovnávací paměť s alespoň požadovaným množstvím paměti. Nepředpokládáme přesné velikosti vyrovnávací paměti.
- Při následných voláních není zaručeno, že vrátí stejnou vyrovnávací paměť nebo vyrovnávací paměť stejné velikosti.
- Po volání Advance musí být požadována nová vyrovnávací paměť, aby bylo možné pokračovat v zápisu dalších dat. Do dříve získané vyrovnávací paměti nelze zapisovat.
- Volání GetMemory nebo GetSpan není bezpečné, pokud probíhá nedokončený hovor FlushAsync.
- Volání Complete nebo CompleteAsync v době, kdy existují nevyprázdněná data, může mít za následek poškození paměti.
Tipy pro PipeReader a PipeWriter
Pomocí těchto tipů můžete úspěšně používat System.IO.Pipelines třídy:
- Vždy dokončete PipeReader a PipeWriter, včetně výjimky, pokud je to možné.
- Vždy volejte PipeReader.AdvanceTo po zavolání PipeReader.ReadAsync.
-
awaitPipeWriter.FlushAsync Pravidelně při psaní a vždy kontrolovat FlushResult.IsCompleted. Přerušte psaní, pokudIsCompletedjetrue, protože to znamená, že čtenář je dokončen a už se nezajímá o to, co je napsané. - Zavolejte PipeWriter.FlushAsync po napsání něčeho, k čemu má mít
PipeReaderpřístup. - Nezavolejte
FlushAsync, pokud čtečka nemůže začít, dokud seFlushAsyncnedokončí, protože to může způsobit deadlock. - Ujistěte se, že pouze jeden kontext vlastní
PipeReaderneboPipeWriter, nebo k nim přistupuje. Tyto typy nejsou vláknotěsné. - Nikdy nepřistupujte k ReadResult.Buffer po volání PipeReader.AdvanceTo nebo dokončení
PipeReader.
IDuplexPipe
IDuplexPipe je kontrakt pro typy, které podporují čtení i psaní. Například síťové připojení je reprezentováno IDuplexPipe.
Na rozdíl od Pipe, který obsahuje PipeReader a a PipeWriter, IDuplexPipe představuje jednu stranu úplného duplexního připojení. To, co píšete do PipeWriter, nebude čteno z PipeReader.
Streamy
Při čtení nebo zápisu dat datového proudu obvykle čtete data pomocí de-serializátoru a zapisujete data pomocí serializátoru. Většina těchto rozhraní API pro čtení a zápis streamu má Stream parametr. Pro usnadnění integrace s těmito stávajícími API PipeReader a PipeWriter nabízejí metodu AsStream.
AsStream vrátí implementaci okolo Stream nebo PipeReader.
Příklad streamu
Vytvářejte PipeReader a PipeWriter instance pomocí statických Create metod zadaných objektem Stream a volitelnými odpovídajícími možnostmi vytvoření.
StreamPipeReaderOptions umožňují kontrolu nad vytvořením instance PipeReader pomocí následujících parametrů:
-
StreamPipeReaderOptions.BufferSize je minimální velikost vyrovnávací paměti v bajtech používaná při pronájmu paměti z fondu a výchozí hodnota je
4096. -
StreamPipeReaderOptions.LeaveOpen příznak určuje, zda po dokončení
PipeReaderzůstane podkladový datový proud otevřený nebo ne, a výchozí hodnota jefalse. -
StreamPipeReaderOptions.MinimumReadSize představuje prahovou hodnotu zbývajících bajtů ve vyrovnávací paměti před přidělením nové vyrovnávací paměti a výchozí hodnota
1024je . -
StreamPipeReaderOptions.Pool
MemoryPool<byte>je použit při přidělování paměti a výchozí hodnota jenull.
StreamPipeWriterOptions umožňují kontrolu nad vytvořením instance PipeWriter pomocí následujících parametrů:
-
StreamPipeWriterOptions.LeaveOpen příznak určuje, zda po dokončení
PipeWriterzůstane podkladový datový proud otevřený nebo ne, a výchozí hodnota jefalse. -
StreamPipeWriterOptions.MinimumBufferSize představuje minimální velikost vyrovnávací paměti, která se má použít při pronájmu paměti z Pool, a výchozí hodnotou je
4096. -
StreamPipeWriterOptions.Pool
MemoryPool<byte>je použit při přidělování paměti a výchozí hodnota jenull.
Důležité
Při vytváření PipeReader a PipeWriter instancí pomocí Create metod zvažte dobu života objektu Stream . Pokud potřebujete přístup ke streamu po dokončení čtení nebo zápisu, nastavte LeaveOpen příznak na true možnosti vytvoření. V opačném případě se proud zavře.
Tento kód ukazuje vytváření PipeReader a PipeWriter instance pomocí Create metod z datového proudu.
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));
}
Aplikace používá StreamReader ke čtení lorem-ipsum.txt souboru jako streamu a musí končit prázdným řádkem. Objekt FileStream je předán PipeReader.Create, který vytvoří instanci objektu PipeReader. Konzolová aplikace pak předá standardní výstupní datový proud na PipeWriter.Create použitím Console.OpenStandardOutput(). Příklad podporuje zrušení.