Асинхронное чтение одного Stream двумя потребителями

Пытался придумать, как одновременно писать получаемые с сервера данные в десериализтор и файл, используемый в будущем как кеш для повторяющихся запросов. Ищу решение, чтобы не читать полностью ответ с сервера и не копировать его вручную побуферно.

Вот моя жалкая попытка, закомментировал то, что у меня не получилось. С закомментированными строчками все работает, а если раскомментировать - десериализатор падает, по вполне понятным причинам.

private static async Task<T> HttpAPIRequest<T>(string url, string path)
{
    using HttpResponseMessage response = await client.GetAsync(url, HttpCompletionOption.ResponseHeadersRead).ConfigureAwait(false);
    response.EnsureSuccessStatusCode();
    using Stream responseStream = await response.Content.ReadAsStreamAsync().ConfigureAwait(false);
    //using FileStream fileStream = File.Create(path);
    //Task fileTask = responseStream.CopyToAsync(fileStream);
    ValueTask<T> jsonTask = JsonSerializer.DeserializeAsync<T>(responseStream);
    //await fileTask;
    return await jsonTask;
}

Но как такое сделать правильно?

У меня есть сейчас рабочий код, но там сначала я полностью гружу ответ в byte[] массив, а затем отправляю этот массив в десериализатор и файл.


Ответы (2 шт):

Автор решения: Alexander Petrov

Как насчёт такого способа:

private static async Task<T> HttpAPIRequest<T>(string url, string path)
{
    using HttpResponseMessage response = await client.GetAsync(url, HttpCompletionOption.ResponseHeadersRead).ConfigureAwait(false);
    response.EnsureSuccessStatusCode();

    using Stream responseStream = await response.Content.ReadAsStreamAsync().ConfigureAwait(false);
    using FileStream fileStream = File.OpenWrite(path);
    using MemoryStream memoryStream = new MemoryStream();

    byte[] buffer = new byte[4096]; // выбрать размер буфера
    int bytesRead;
    while ((bytesRead = await responseStream.ReadAsync(buffer, 0, buffer.Length)) > 0)
    {
        var task1 = fileStream.WriteAsync(buffer, 0, bytesRead);
        var task2 = memoryStream.WriteAsync(buffer, 0, bytesRead);
        await Task.WhenAll(task1, task2);
    }
    // memoryStream.Position = 0;
    ValueTask<T> jsonTask = JsonSerializer.DeserializeAsync<T>(memoryStream);
    return await jsonTask;
}
→ Ссылка
Автор решения: aepot

Пока я ждал ответа, прогресс не стоял на месте, и у меня все-таки удалось раскусить этот вопрос.

Суть ТЗ:

  • Форкнуть стрим асинхронно так, чтобы он одновременно писался в несколько стримов-приемников
  • При этом, каждый из приемников может завершить свою работу раньше остальных
  • Чтение из стрима-источника должно происходить с такой скоростью, с какой читает самый быстрый приемник
  • Если поток-источник медленный (а он всегда медленный в моем случае), то приемники должны ожидать появления новых данных

Решение пришло откуда не ждали: IEnumerator.

public sealed class ConcurrentStreamReader : IDisposable
{
    private readonly Stream _inputStream;
    private readonly MemoryStream _memoryStream;
    private readonly int _bufferSize;
    private readonly IProgress<long> _progress;
    private bool disposed;

    public ConcurrentStreamReader(Stream inputStream) : this(inputStream, null) { }

    public ConcurrentStreamReader(Stream inputStream, IProgress<long> progress) : this(inputStream, 81920, 0, progress) { }

    public ConcurrentStreamReader(Stream inputStream, int bufferSize, int capacity, IProgress<long> progress)
    {
        _memoryStream = new MemoryStream(capacity);
        _inputStream = inputStream;
        _bufferSize = bufferSize;
        _progress = progress;
    }

    public Stream CreateNewStream() => new BlockingStream(ReadStream().GetEnumerator(), _memoryStream);

    public async Task CopyToAsync(Stream outputStream)
    {
        foreach (ReadOnlyMemory<byte> buffer in ReadStream())
        {
            await outputStream.WriteAsync(buffer);
        }
    }

    private IEnumerable<ReadOnlyMemory<byte>> ReadStream()
    {
        long position = 0;
        byte[] buffer = new byte[_bufferSize];
        int bytesRead;
        while (true)
        {
            if (position == _memoryStream.Length)
            {
                lock (_inputStream)
                {
                    if (position < _memoryStream.Length) continue;
                    bytesRead = _inputStream.Read(buffer, 0, _bufferSize);
                    if (bytesRead == 0) break;
                    lock (_memoryStream)
                    {
                        _memoryStream.Position = _memoryStream.Length;
                        _memoryStream.Write(buffer, 0, bytesRead);
                    }
                }
                _progress?.Report(_memoryStream.Length);
            }
            else
            {
                lock (_memoryStream)
                {
                    _memoryStream.Position = position;
                    bytesRead = _memoryStream.Read(buffer, 0, _bufferSize);
                }
            }
            position += bytesRead;
            yield return new ReadOnlyMemory<byte>(buffer, 0, bytesRead);
        }
    }

    public void Dispose()
    {
        Dispose(true);
        GC.SuppressFinalize(this);
    }
    private void Dispose(bool disposing)
    {
        if (disposing && !disposed)
        {
            disposed = true;
            _memoryStream?.Dispose();
        }
    }
}

public sealed class BlockingStream : Stream
{
    private readonly IEnumerator<ReadOnlyMemory<byte>> _sourceEnumerator;
    private readonly MemoryStream _stream;
    private long _position;

    private BlockingStream() { }

    public BlockingStream(IEnumerator<ReadOnlyMemory<byte>> sourceEnumerator, MemoryStream stream)
    {
        _stream = stream;
        _sourceEnumerator = sourceEnumerator;
        _position = 0;
    }

    public override int Read(byte[] buffer, int offset, int count)
    {
        long tail = _position + count;
        while (tail >= _stream.Length && _sourceEnumerator.MoveNext()) { }
        int bytesRead;
        lock (_stream)
        {
            _stream.Position = _position;
            bytesRead = _stream.Read(buffer, offset, count);
        }
        _position += bytesRead;
        return bytesRead;
    }

    public override bool CanRead => true;
    public override bool CanSeek => false;
    public override bool CanWrite => false;
    public override long Length => throw new NotSupportedException();
    public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }
    public override void Flush() => throw new NotSupportedException();
    public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
    public override void SetLength(long value) => throw new NotSupportedException();
    public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException();
}

И вот так вот оно теперь используется.

private async Task<T> HttpAPIRequest<T>(string url, string path)
{
    using HttpResponseMessage response = await _client.GetAsync(url, HttpCompletionOption.ResponseHeadersRead).ConfigureAwait(false);
    response.EnsureSuccessStatusCode();
    using Stream responseStream = await response.Content.ReadAsStreamAsync().ConfigureAwait(false);
    using ConcurrentStreamReader reader = new ConcurrentStreamReader(responseStream);
    using FileStream fileStream = File.Create(path);
    using Stream deserealizerStream = reader.CreateNewStream();
    ValueTask<T> jsonTask = JsonSerializer.DeserializeAsync<T>(deserealizerStream);
    await reader.CopyToAsync(fileStream);
    return await jsonTask;
}

Если есть советы, что допилить, или где мне не хватило ума сделать нормально - пишите в комментариях. Всем спасибо, кто пытался помочь.

→ Ссылка