欢迎您访问程序员文章站本站旨在为大家提供分享程序员计算机编程知识!
您现在的位置是: 首页  >  IT编程

Pipe——高性能IO

程序员文章站 2023-11-15 17:12:34
System.IO.Pipelines是一个新的库,旨在简化在.NET中执行高性能IO的过程。它是一个依赖.NET Standard的库,适用于所有.NET实现。 Pipelines诞生于.NET Core团队,为使Kestrel成为业界最快的Web服务器之一。最初从作为Kestrel内部的实现细节 ......

system.io.pipelines是一个新的库,旨在简化在.net中执行高性能io的过程。它是一个依赖.net standard的库,适用于所有.net实现

pipelines诞生于.net core团队,为使kestrel成为业界最快的web服务器之一。最初从作为kestrel内部的实现细节发展成为可重用的api,它在.net core 2.1中作为可用于所有.net开发人员的*bcl api(system.io.pipelines)提供。

它解决了什么问题?

为了正确解析stream或socket中的数据,代码有固定的样板,并且有许多极端情况,为了处理他们,不得不编写难以维护的复杂代码。
实现高性能和正确性,同时也难以处理这种复杂性。pipelines旨在解决这种复杂性。

有多复杂?

让我们从一个简单的问题开始吧。我们想编写一个tcp服务器,它接收来自客户端的用行分隔的消息(由\n分隔)。(译者注:即一行为一条消息)

使用networkstream的tcp服务器

在pipelines之前用.net编写的典型代码如下所示:

async task processlinesasync(networkstream stream)
{
    var buffer = new byte[1024];
    await stream.readasync(buffer, 0, buffer.length);
    
    // 在buffer中处理一行消息
    processline(buffer);
}

此代码可能在本地测试时正确工作,但它有几个潜在错误:

  • 一次readasync调用可能没有收到整个消息(行尾)。
  • 它忽略了stream.readasync()返回值中实际填充到buffer中的数据量。(译者注:即不一定将buffer填充满)
  • 一次readasync调用不能处理多条消息。

这些是读取流数据时常见的一些缺陷。为了解决这个问题,我们需要做一些改变:

  • 我们需要缓冲传入的数据,直到找到新的行。
  • 我们需要解析缓冲区中返回的所有行
async task processlinesasync(networkstream stream)
{
    var buffer = new byte[1024];
    var bytesbuffered = 0;
    var bytesconsumed = 0;

    while (true)
    {
        var bytesread = await stream.readasync(buffer, bytesbuffered, buffer.length - bytesbuffered);
        if (bytesread == 0)
        {
            // eof 已经到末尾
            break;
        }
        // 跟踪已缓冲的字节数
        bytesbuffered += bytesread;
        
        var lineposition = -1;

        do
        {
            // 在缓冲数据中查找找一个行末尾
            lineposition = array.indexof(buffer, (byte)‘\n‘, bytesconsumed, bytesbuffered - bytesconsumed);

            if (lineposition >= 0)
            {
                // 根据偏移量计算一行的长度
                var linelength = lineposition - bytesconsumed;

                // 处理这一行
                processline(buffer, bytesconsumed, linelength);

                // 移动bytesconsumed为了跳过我们已经处理掉的行 (包括\n)
                bytesconsumed += linelength + 1;
            }
        }
        while (lineposition >= 0);
    }
}

这一次,这可能适用于本地开发,但一行可能大于1kib(1024字节)。我们需要调整输入缓冲区的大小,直到找到新行。

因此,我们可以在堆上分配缓冲区去处理更长的一行。我们从客户端解析较长的一行时,可以通过使用arraypool<byte>避免重复分配缓冲区来改进这一点。

async task processlinesasync(networkstream stream)
{
    byte[] buffer = arraypool<byte>.shared.rent(1024);
    var bytesbuffered = 0;
    var bytesconsumed = 0;

    while (true)
    {
        // 在buffer中计算中剩余的字节数
        var bytesremaining = buffer.length - bytesbuffered;

        if (bytesremaining == 0)
        {
            // 将buffer size翻倍 并且将之前缓冲的数据复制到新的缓冲区
            var newbuffer = arraypool<byte>.shared.rent(buffer.length * 2);
            buffer.blockcopy(buffer, 0, newbuffer, 0, buffer.length);
            // 将旧的buffer丢回池中
            arraypool<byte>.shared.return(buffer);
            buffer = newbuffer;
            bytesremaining = buffer.length - bytesbuffered;
        }

        var bytesread = await stream.readasync(buffer, bytesbuffered, bytesremaining);
        if (bytesread == 0)
        {
            // eof 末尾
            break;
        }
        
        // 跟踪已缓冲的字节数
        bytesbuffered += bytesread;
        
        do
        {
            // 在缓冲数据中查找找一个行末尾
            lineposition = array.indexof(buffer, (byte)‘\n‘, bytesconsumed, bytesbuffered - bytesconsumed);

            if (lineposition >= 0)
            {
                // 根据偏移量计算一行的长度
                var linelength = lineposition - bytesconsumed;

                // 处理这一行
                processline(buffer, bytesconsumed, linelength);

                // 移动bytesconsumed为了跳过我们已经处理掉的行 (包括\n)
                bytesconsumed += linelength + 1;
            }
        }
        while (lineposition >= 0);
    }
}

这段代码有效,但现在我们正在重新调整缓冲区大小,从而产生更多缓冲区副本。它将使用更多内存,因为根据代码在处理一行行后不会缩缓冲区的大小。为避免这种情况,我们可以存储缓冲区序列,而不是每次超过1kib大小时调整大小。

此外,我们不会增长1kib的 缓冲区,直到它完全为空。这意味着我们最终传递给readasync越来越小的缓冲区,这将导致对操作系统的更多调用。

为了缓解这种情况,我们将在现有缓冲区中剩余少于512个字节时分配一个新缓冲区:

public class buffersegment
{
    public byte[] buffer { get; set; }
    public int count { get; set; }

    public int remaining => buffer.length - count;
}

async task processlinesasync(networkstream stream)
{
    const int minimumbuffersize = 512;

    var segments = new list<buffersegment>();
    var bytesconsumed = 0;
    var bytesconsumedbufferindex = 0;
    var segment = new buffersegment { buffer = arraypool<byte>.shared.rent(1024) };

    segments.add(segment);

    while (true)
    {
        // calculate the amount of bytes remaining in the buffer
        if (segment.remaining < minimumbuffersize)
        {
            // allocate a new segment
            segment = new buffersegment { buffer = arraypool<byte>.shared.rent(1024) };
            segments.add(segment);
        }

        var bytesread = await stream.readasync(segment.buffer, segment.count, segment.remaining);
        if (bytesread == 0)
        {
            break;
        }

        // keep track of the amount of buffered bytes
        segment.count += bytesread;

        while (true)
        {
            // look for a eol in the list of segments
            var (segmentindex, segmentoffset) = indexof(segments, (byte)‘\n‘, bytesconsumedbufferindex, bytesconsumed);

            if (segmentindex >= 0)
            {
                // process the line
                processline(segments, segmentindex, segmentoffset);

                bytesconsumedbufferindex = segmentoffset;
                bytesconsumed = segmentoffset + 1;
            }
            else
            {
                break;
            }
        }

        // drop fully consumed segments from the list so we don‘t look at them again
        for (var i = bytesconsumedbufferindex; i >= 0; --i)
        {
            var consumedsegment = segments[i];
            // return all segments unless this is the current segment
            if (consumedsegment != segment)
            {
                arraypool<byte>.shared.return(consumedsegment.buffer);
                segments.removeat(i);
            }
        }
    }
}

(int segmentindex, int segmentoffest) indexof(list<buffersegment> segments, byte value, int startbufferindex, int startsegmentoffset)
{
    var first = true;
    for (var i = startbufferindex; i < segments.count; ++i)
    {
        var segment = segments[i];
        // start from the correct offset
        var offset = first ? startsegmentoffset : 0;
        var index = array.indexof(segment.buffer, value, offset, segment.count - offset);

        if (index >= 0)
        {
            // return the buffer index and the index within that segment where eol was found
            return (i, index);
        }

        first = false;
    }
    return (-1, -1);
}

此代码只是得到很多更加复杂。当我们正在寻找分隔符时,我们同时跟踪已填充的缓冲区序列。为此,我们此处使用list<buffersegment>查找新行分隔符时表示缓冲数据。其结果是,processlineindexof现在接受list<buffersegment>作为参数,而不是一个byte[],offset和count。我们的解析逻辑现在需要处理一个或多个缓冲区序列。

我们的服务器现在处理部分消息,它使用池化内存来减少总体内存消耗,但我们还需要进行更多更改:

  1. 我们使用的byte[]arraypool<byte>的只是普通的托管数组。这意味着无论何时我们执行readasyncwriteasync,这些缓冲区都会在异步操作的生命周期内被固定(以便与操作系统上的本机io api互操作)。这对gc有性能影响,因为无法移动固定内存,这可能导致堆碎片。根据异步操作挂起的时间长短,池的实现可能需要更改。
  2. 可以通过解耦读取逻辑处理逻辑来优化吞吐量。这会创建一个批处理效果,使解析逻辑可以使用更大的缓冲区块,而不是仅在解析单个行后才读取更多数据。这引入了一些额外的复杂性
    • 我们需要两个彼此独立运行的循环。一个读取socket和一个解析缓冲区。
    • 当数据可用时,我们需要一种方法来向解析逻辑发出信号。
    • 我们需要决定如果循环读取socket“太快”会发生什么。如果解析逻辑无法跟上,我们需要一种方法来限制读取循环(逻辑)。这通常被称为“流量控制”或“背压”。
    • 我们需要确保事情是线程安全的。我们现在在读取循环解析循环之间共享多个缓冲区,并且这些缓冲区在不同的线程上独立运行。
    • 内存管理逻辑现在分布在两个不同的代码段中,从填充缓冲区池的代码是从套接字读取的,而从缓冲区池取数据的代码是解析逻辑
    • 我们需要非常小心在解析逻辑完成之后我们如何处理缓冲区序列。如果我们不小心,我们可能会返回一个仍由socket读取逻辑写入的缓冲区序列。

复杂性已经到了极端(我们甚至没有涵盖所有案例)。高性能网络应用通常意味着编写非常复杂的代码,以便从系统中获得更高的性能。

system.io.pipelines的目标是使这种类型的代码更容易编写。

使用system.io.pipelines的tcp服务器

让我们来看看这个例子的样子system.io.pipelines:

async task processlinesasync(socket socket)
{
    var pipe = new pipe();
    task writing = fillpipeasync(socket, pipe.writer);
    task reading = readpipeasync(pipe.reader);

    return task.whenall(reading, writing);
}

async task fillpipeasync(socket socket, pipewriter writer)
{
    const int minimumbuffersize = 512;

    while (true)
    {
        // 从pipewriter至少分配512字节
        memory<byte> memory = writer.getmemory(minimumbuffersize);
        try 
        {
            int bytesread = await socket.receiveasync(memory, socketflags.none);
            if (bytesread == 0)
            {
                break;
            }
            // 告诉pipewriter从套接字读取了多少
            writer.advance(bytesread);
        }
        catch (exception ex)
        {
            logerror(ex);
            break;
        }

        // 标记数据可用,让pipereader读取
        flushresult result = await writer.flushasync();

        if (result.iscompleted)
        {
            break;
        }
    }

    // 告诉pipereader没有更多的数据
    writer.complete();
}

async task readpipeasync(pipereader reader)
{
    while (true)
    {
        readresult result = await reader.readasync();

        readonlysequence<byte> buffer = result.buffer;
        sequenceposition? position = null;

        do 
        {
            // 在缓冲数据中查找找一个行末尾
            position = buffer.positionof((byte)‘\n‘);

            if (position != null)
            {
                // 处理这一行
                processline(buffer.slice(0, position.value));
                
                // 跳过 这一行+\n (basically position 主要位置?)
                buffer = buffer.slice(buffer.getposition(1, position.value));
            }
        }
        while (position != null);

        // 告诉pipereader我们以及处理多少缓冲
        reader.advanceto(buffer.start, buffer.end);

        // 如果没有更多的数据,停止都去
        if (result.iscompleted)
        {
            break;
        }
    }

    // 将pipereader标记为完成
    reader.complete();
}

我们的行读取器的pipelines版本有2个循环:

  • fillpipeasync从socket读取并写入pipewriter。
  • readpipeasync从pipereader中读取并解析传入的行。

与原始示例不同,在任何地方都没有分配显式缓冲区。这是管道的核心功能之一。所有缓冲区管理都委托给pipereader/pipewriter实现。

这使得使用代码更容易专注于业务逻辑而不是复杂的缓冲区管理。

在第一个循环中,我们首先调用pipewriter.getmemory(int)从底层编写器获取一些内存; 然后我们调用pipewriter.advance(int)告诉pipewriter我们实际写入缓冲区的数据量。然后我们调用pipewriter.flushasync()来提供数据给pipereader。

在第二个循环中,我们正在使用pipewriter最终来自的缓冲区socket。当调用pipereader.readasync()返回时,我们得到一个readresult包含2条重要信息,包括以readonlysequence<byte>形式读取的数据和bool iscompleted,让reader知道writer是否写完(eof)。在找到行尾(eol)分隔符并解析该行之后,我们将缓冲区切片以跳过我们已经处理过的内容,然后我们调用pipereader.advanceto告诉pipereader我们消耗了多少数据。

在每个循环结束时,我们完成了reader和writer。这允许底层pipe释放它分配的所有内存。

system.io.pipelines

除了处理内存管理之外,其他核心管道功能还包括能够在pipe不实际消耗数据的情况下查看数据。

pipereader有两个核心api readasyncadvancetoreadasync获取pipe数据,advanceto告诉pipereader不再需要这些缓冲区,以便可以丢弃它们(例如返回到底层缓冲池)。


这是一个http解析器的示例,它在接收pipe到有效起始行之前读取部分数据缓冲区数据。

Pipe——高性能IO

readonlysequence<t>

该pipe实现存储了在pipewriter和pipereader之间传递的缓冲区的链接列表。pipereader.readasync暴露一个readonlysequence<t>新的bcl类型,它表示一个或多个readonlymemory<t>段的视图,类似于span<t>和memory<t>提供数组和字符串的视图。

Pipe——高性能IO

该pipe内部维护指向reader和writer可以分配或更新它们的数据集合,。sequenceposition表示缓冲区链表中的单个点,可用于有效地对readonlysequence<t>进行切片。

这段实在翻译困难,给出原文
the pipe internally maintains pointers to where the reader and writer are in the overall set of allocated data and updates them as data is written or read. the sequenceposition represents a single point in the linked list of buffers and can be used to efficiently slice the readonlysequence

由于readonlysequence<t>可以支持一个或多个段,因此高性能处理逻辑通常基于单个或多个段来分割快速和慢速路径(fast and slow paths?)。

例如,这是一个将ascii readonlysequence<byte>转换为string以下内容的例程:

string getasciistring(readonlysequence<byte> buffer)
{
    if (buffer.issinglesegment)
    {
        return encoding.ascii.getstring(buffer.first.span);
    }

    return string.create((int)buffer.length, buffer, (span, sequence) =>
    {
        foreach (var segment in sequence)
        {
            encoding.ascii.getchars(segment.span, span);

            span = span.slice(segment.length);
        }
    });
}

背压和流量控制

在一个完美的世界中,读取和解析工作是一个团队:读取线程消耗来自网络的数据并将其放入缓冲区,而解析线程负责构建适当的数据结构。通常,解析将比仅从网络复制数据块花费更多时间。结果,读取线程可以轻易地压倒解析线程。结果是读取线程必须减慢或分配更多内存来存储解析线程的数据。为获得最佳性能,在频繁暂停和分配更多内存之间存在平衡。

为了解决这个问题,管道有两个设置来控制数据的流量,pausewriterthreshold和resumewriterthreshold。pausewriterthreshold决定有多少数据应该在调用pipewriter.flushasync之前进行缓冲停顿。resumewriterthreshold控制reader消耗多少后写入可以恢复。

Pipe——高性能IO

当pipe的数据量超过pausewriterthreshold,pipewriter.flushasync会异步阻塞。数据量变得低于resumewriterthreshold,它会解锁时。两个值用于防止在极限附近发生反复阻塞和解锁。

io调度

通常在使用async / await时,会在线程池线程或当前线程上调用continuation synchronizationcontext。

在执行io时,对执行io的位置进行细粒度控制非常重要,这样可以更有效地利用cpu缓存,这对于web服务器等高性能应用程序至关重要。pipelines公开了一个pipescheduler确定异步回调运行位置的方法。这使得调用者可以精确控制用于io的线程。

实践中的一个示例是在kestrel libuv传输中,其中io回调在专用事件循环线程上运行。

pipereader模式的其他好处:

  • 一些底层系统支持“无缓冲等待”,即,在底层系统中实际可用数据之前,永远不需要分配缓冲区。例如,在带有epoll的linux上,可以等到数据准备好之后再实际提供缓冲区来进行读取。这避免了具有大量线程等待数据的问题不会立即需要保留大量内存。
  • 默认情况下pipe,可以轻松地针对网络代码编写单元测试,因为解析逻辑与网络代码分离,因此单元测试仅针对内存缓冲区运行解析逻辑,而不是直接从网络中消耗。它还可以轻松测试那些难以测试发送部分数据的模式。asp.net core使用它来测试kestrel的http解析器的各个方面。
  • 允许将底层os缓冲区(如windows上的registered io api)暴露给用户代码的系统非常适合管道,因为缓冲区始终由pipereader实现提供。

其他相关类型

作为制作system.io.pipelines的一部分,我们还添加了许多新的原始bcl类型:

  • memorypool<t>imemoryowner<t>memorymanager<t> - .net core 1.0添加了arraypool<t>,在.net core 2.1中,我们现在有一个更通用的抽象,适用于任何工作的池memory<t>。这提供了一个可扩展点,允许您插入更高级的分配策略以及控制缓冲区的管理方式(例如,提供预先固定的缓冲区而不是纯托管的阵列)。
  • ibufferwriter<t> - 表示用于写入同步缓冲数据的接收器。(pipewriter实现这个)
  • ivaluetasksource - valuetask<t>自.net core 1.1以来就已存在,但在.net core 2.1中获得了一些超级权限,允许无分配的等待异步操作。有关详细信息,请参阅https://github.com/dotnet/corefx/issues/27445。

如何使用管道?

api存在于system.io.pipelines nuget包中。

主要包含一个pipe对象,它有一个writer属性和reader属性。

var pipe = new pipe();
var writer = pipe.writer;
var reader = pipe.reader; 

writer对象

writer对象用于从数据源读取数据,将数据写入管道中;它对应业务中的"读"操作。

var content = encoding.default.getbytes("hello world");
var data = new memory<byte>(content);
var result = await writer.writeasync(data);

另外,它也有一种使用pipe申请memory的方式

var buffer = writer.getmemory(512);
content.copyto(buffer);
writer.advance(content.length);
var result = await writer.flushasync(); 

reader对象

reader对象用于从管道中获取数据源,它对应业务中的"用"操作。

首先获取管道的缓冲区:

var result = await reader.readasync();
var buffer = result.buffer;

这个buffer是一个readonlysequence<byte>对象,它是一个相当好的动态内存对象,并且相当高效。它本身由多段memory<byte>组成,查看memory段的方法有:

issinglesegment: 判断是否只有一段memory<byte>
first: 获取第一段memory<byte>
getenumerator: 获取分段的memory<byte>
它从逻辑上也可以看成一段连续的memory<byte>,也有类似的方法:

length: 整个数据缓冲区长度
slice: 分割缓冲区
copyto: 将内容复制到span中
toarray: 将内容复制到byte[]中
另外,它还有一个类似游标的位置对象sequenceposition,可以从其position相关函数中使用,这里就不多介绍了。

这个缓冲区解决了"数据读不够"的问题,一次读取的不够下次可以接着读,不用缓冲区的动态分配,高效的内存管理方式带来了良好的性能,好用的接口是我们能更关注业务。

获取到缓冲区后,就是使用缓冲区的数据

var data = buffer.toarray();

使用完后,告诉pipe当前使用了多少数据,下次接着从结束位置后读起

reader.advanceto(buffer.getposition(4));

这是一个相当实用的设计,它解决了"读了就得用"的问题,不仅可以将不用的数据下次再使用,还可以实现peek的操作,只读但不改变游标。 

交互

除了"读"和"用"操作外,它们之间还需要一些交互,例如:

读过程中数据源不可用,需要停止使用
使用过程中业务结束,需要中止数据源。
reader和writer都有一个complete函数,用于通知结束:

reader.complete();
writer.complete();

在writer写入和reader读取时,会获得一个结果

flushresult result = await writer.flushasync();
readresult result = await reader.readasync();

它们都有一个iscomplete属性,可以根据它是否为true判断是否已经结束了读和写的操作。 

取消

在写入和读取的时候,也可以传入一个cancellationtoken,用于取消相应的操作。

writer.flushasync(cancellationtoken.none);
reader.readasync(cancellationtoken.none);

如果取消成功,对应的result的iscanceled则为true

 


转载请标明本文来源:
更多内容欢迎star、fork我的的github:
如果发现本文有什么问题和任何建议,也随时欢迎交流~