Nguon: Microsoft Learn · .NET 8.0

Sử dụng Streaming trong ASP.NET Core SignalR

Nguồn: Use streaming in ASP.NET Core SignalR

Tác giả: Brennan Conroy

ASP.NET Core SignalR hỗ trợ streaming (truyền dữ liệu liên tục) từ client đến server và từ server đến client. Điều này hữu ích cho các tình huống mà các mảnh dữ liệu đến theo thời gian. Khi streaming, mỗi mảnh được gửi đến client hoặc server ngay khi nó có sẵn, thay vì chờ đợi tất cả dữ liệu trở nên sẵn sàng.

Xem hoặc tải xuống mã mẫu

Thiết lập hub cho streaming

Một phương thức hub tự động trở thành phương thức hub streaming khi nó trả về IAsyncEnumerable<T>, ChannelReader<T>, Task<IAsyncEnumerable<T>> hoặc Task<ChannelReader<T>>.

Streaming từ server đến client (Server-to-client streaming)

Các phương thức hub streaming có thể trả về IAsyncEnumerable<T> ngoài ChannelReader<T>. Cách đơn giản nhất để trả về IAsyncEnumerable<T> là biến phương thức hub thành phương thức async iterator như ví dụ sau. Các phương thức hub async iterator có thể chấp nhận tham số CancellationToken được kích hoạt khi client hủy đăng ký stream. Các phương thức async iterator tránh các vấn đề phổ biến với Channels, chẳng hạn như không trả về ChannelReader đủ sớm hoặc thoát khỏi phương thức mà không hoàn thành ChannelWriter<T>.

Lưu ý: Ví dụ sau yêu cầu C# 8.0 trở lên.

csharp
public class AsyncEnumerableHub : Hub
{
    public async IAsyncEnumerable<int> Counter(
        int count,
        int delay,
        [EnumeratorCancellation]
        CancellationToken cancellationToken)
    {
        for (var i = 0; i < count; i++)
        {
            // Kiểm tra cancellation token thường xuyên để server dừng
            // tạo ra các item nếu client ngắt kết nối.
            cancellationToken.ThrowIfCancellationRequested();

            yield return i;

            // Sử dụng cancellationToken trong các API khác chấp nhận cancellation
            // tokens để cancellation có thể truyền xuống chúng.
            await Task.Delay(delay, cancellationToken);
        }
    }
}

Ví dụ sau hiển thị cơ bản về streaming dữ liệu đến client sử dụng Channels. Bất cứ khi nào một đối tượng được ghi vào ChannelWriter<T>, đối tượng đó ngay lập tức được gửi đến client. Cuối cùng, ChannelWriter được hoàn thành để thông báo cho client rằng stream đã đóng.

Lưu ý: Ghi vào ChannelWriter<T> trên background thread và trả về ChannelReader càng sớm càng tốt. Các lời gọi hub khác bị chặn cho đến khi ChannelReader được trả về.

Bọc logic trong một câu lệnh try...catch. Hoàn thành Channel trong khối finally. Nếu bạn muốn truyền lỗi, hãy bắt nó trong khối catch và ghi vào khối finally.

csharp
public ChannelReader<int> Counter(
    int count,
    int delay,
    CancellationToken cancellationToken)
{
    var channel = Channel.CreateUnbounded<int>();

    // Chúng ta không muốn await WriteItemsAsync, nếu không chúng ta sẽ phải
    // chờ tất cả các item được ghi trước khi trả về channel về cho client.
    _ = WriteItemsAsync(channel.Writer, count, delay, cancellationToken);

    return channel.Reader;
}

private async Task WriteItemsAsync(
    ChannelWriter<int> writer,
    int count,
    int delay,
    CancellationToken cancellationToken)
{
    Exception localException = null;
    try
    {
        for (var i = 0; i < count; i++)
        {
            await writer.WriteAsync(i, cancellationToken);

            // Sử dụng cancellationToken trong các API khác chấp nhận cancellation
            // tokens để cancellation có thể truyền xuống chúng.
            await Task.Delay(delay, cancellationToken);
        }
    }
    catch (Exception ex)
    {
        localException = ex;
    }
    finally
    {
        writer.Complete(localException);
    }
}

Các phương thức hub streaming server-to-client có thể chấp nhận tham số CancellationToken được kích hoạt khi client hủy đăng ký stream. Sử dụng token này để dừng hoạt động server và giải phóng bất kỳ tài nguyên nào nếu client ngắt kết nối trước khi stream kết thúc.

Streaming từ client đến server (Client-to-server streaming)

Một phương thức hub tự động trở thành phương thức hub client-to-server streaming khi nó chấp nhận một hoặc nhiều đối tượng kiểu ChannelReader<T> hoặc IAsyncEnumerable<T>. Ví dụ sau hiển thị cơ bản về đọc dữ liệu streaming được gửi từ client. Bất cứ khi nào client ghi vào ChannelWriter<T>, dữ liệu được ghi vào ChannelReader trên server mà phương thức hub đang đọc từ đó.

csharp
public async Task UploadStream(ChannelReader<string> stream)
{
    while (await stream.WaitToReadAsync())
    {
        while (stream.TryRead(out var item))
        {
            // làm gì đó với stream item
            Console.WriteLine(item);
        }
    }
}

Phiên bản IAsyncEnumerable<T> của phương thức như sau.

Lưu ý: Ví dụ sau yêu cầu C# 8.0 trở lên.

csharp
public async Task UploadStream(IAsyncEnumerable<string> stream)
{
    await foreach (var item in stream)
    {
        Console.WriteLine(item);
    }
}

.NET client

Streaming từ server đến client

Các phương thức StreamAsyncStreamAsChannelAsync trên HubConnection được sử dụng để gọi các phương thức streaming server-to-client. Truyền tên phương thức hub và các đối số được định nghĩa trong phương thức hub vào StreamAsync hoặc StreamAsChannelAsync. Tham số generic trên StreamAsync<T>StreamAsChannelAsync<T> chỉ định kiểu đối tượng được trả về bởi phương thức streaming. Một đối tượng kiểu IAsyncEnumerable<T> hoặc ChannelReader<T> được trả về từ lời gọi stream và đại diện cho stream trên client.

Ví dụ StreamAsync trả về IAsyncEnumerable<int>:

csharp
// Gọi "Cancel" trên CancellationTokenSource này để gửi thông điệp hủy bỏ đến
// server, sẽ kích hoạt token tương ứng trong phương thức hub.
var cancellationTokenSource = new CancellationTokenSource();
var stream = hubConnection.StreamAsync<int>(
    "Counter", 10, 500, cancellationTokenSource.Token);

await foreach (var count in stream)
{
    Console.WriteLine($"{count}");
}

Console.WriteLine("Streaming completed");

Ví dụ StreamAsChannelAsync tương ứng trả về ChannelReader<int>:

csharp
// Gọi "Cancel" trên CancellationTokenSource này để gửi thông điệp hủy bỏ đến
// server, sẽ kích hoạt token tương ứng trong phương thức hub.
var cancellationTokenSource = new CancellationTokenSource();
var channel = await hubConnection.StreamAsChannelAsync<int>(
    "Counter", 10, 500, cancellationTokenSource.Token);

// Chờ bất đồng bộ cho đến khi dữ liệu sẵn có
while (await channel.WaitToReadAsync())
{
    // Đọc tất cả dữ liệu hiện có đồng bộ, trước khi chờ thêm dữ liệu
    while (channel.TryRead(out var count))
    {
        Console.WriteLine($"{count}");
    }
}

Console.WriteLine("Streaming completed");

Streaming từ client đến server

Có hai cách để gọi phương thức hub client-to-server streaming từ .NET client. Bạn có thể truyền vào IAsyncEnumerable<T> hoặc ChannelReader làm đối số cho SendAsync, InvokeAsync hoặc StreamAsChannelAsync, tùy thuộc vào phương thức hub được gọi.

Bất cứ khi nào dữ liệu được ghi vào đối tượng IAsyncEnumerable hoặc ChannelWriter, phương thức hub trên server nhận một item mới với dữ liệu từ client.

Nếu sử dụng đối tượng IAsyncEnumerable, stream kết thúc sau khi phương thức trả về các stream item thoát.

Lưu ý: Ví dụ sau yêu cầu C# 8.0 trở lên.

csharp
async IAsyncEnumerable<string> clientStreamData()
{
    for (var i = 0; i < 5; i++)
    {
        var data = await FetchSomeData();
        yield return data;
    }
    // Sau khi vòng lặp for hoàn thành và hàm cục bộ thoát, 
    // stream completion sẽ được gửi.
}

await connection.SendAsync("UploadStream", clientStreamData());

Hoặc nếu bạn đang sử dụng ChannelWriter, hãy hoàn thành channel với channel.Writer.Complete():

csharp
var channel = Channel.CreateBounded<string>(10);
await connection.SendAsync("UploadStream", channel.Reader);
await channel.Writer.WriteAsync("some data");
await channel.Writer.WriteAsync("some more data");
channel.Writer.Complete();

JavaScript client

Streaming từ server đến client

JavaScript clients gọi các phương thức streaming server-to-client trên hubs bằng connection.stream. Phương thức stream chấp nhận hai đối số:

connection.stream trả về IStreamResult, chứa phương thức subscribe. Truyền IStreamSubscriber vào subscribe và đặt các callback next, errorcomplete để nhận thông báo từ lời gọi stream.

javascript
connection.stream("Counter", 10, 500)
    .subscribe({
        next: (item) => {
            var li = document.createElement("li");
            li.textContent = item;
            document.getElementById("messagesList").appendChild(li);
        },
        complete: () => {
            var li = document.createElement("li");
            li.textContent = "Stream completed";
            document.getElementById("messagesList").appendChild(li);
        },
        error: (err) => {
            var li = document.createElement("li");
            li.textContent = err;
            document.getElementById("messagesList").appendChild(li);
        },
});

Để kết thúc stream từ client, gọi phương thức dispose trên ISubscription được trả về từ phương thức subscribe. Gọi phương thức này gây ra hủy bỏ tham số CancellationToken của phương thức Hub, nếu bạn đã cung cấp.

Streaming từ client đến server

JavaScript clients gọi các phương thức streaming client-to-server trên hubs bằng cách truyền vào Subject làm đối số cho send, invoke hoặc stream, tùy thuộc vào phương thức hub được gọi. Subject là một lớp trông giống Subject. Ví dụ trong RxJS, bạn có thể sử dụng lớp Subject từ thư viện đó.

javascript
const subject = new signalR.Subject();
yield connection.send("UploadStream", subject);
var iteration = 0;
const intervalHandle = setInterval(() => {
    iteration++;
    subject.next(iteration.toString());
    if (iteration === 10) {
        clearInterval(intervalHandle);
        subject.complete();
    }
}, 500);

Gọi subject.next(item) với một item ghi item vào stream, và phương thức hub nhận item trên server.

Để kết thúc stream, gọi subject.complete().

Java client

Streaming từ server đến client

SignalR Java client sử dụng phương thức stream để gọi các phương thức streaming. stream chấp nhận ba hoặc nhiều đối số hơn:

java
hubConnection.stream(String.class, "ExampleStreamingHubMethod", "Arg1")
    .subscribe(
        (item) -> {/* Định nghĩa onNext handler của bạn ở đây. */ },
        (error) -> {/* Định nghĩa onError handler của bạn ở đây. */},
        () -> {/* Định nghĩa onCompleted handler của bạn ở đây. */});

Phương thức stream trên HubConnection trả về một Observable của kiểu stream item. Phương thức subscribe của kiểu Observable là nơi onNext, onErroronCompleted handlers được định nghĩa.

Streaming từ client đến server

SignalR Java client có thể gọi các phương thức streaming client-to-server trên hubs bằng cách truyền vào một Observable làm đối số cho send, invoke hoặc stream, tùy thuộc vào phương thức hub được gọi.

java
ReplaySubject<String> stream = ReplaySubject.create();
hubConnection.send("UploadStream", stream);
stream.onNext("FirstItem");
stream.onNext("SecondItem");
stream.onComplete();

Gọi stream.onNext(item) với một item ghi item vào stream, và phương thức hub nhận item trên server.

Để kết thúc stream, gọi stream.onComplete().