使用 C# 和 .NET 生成和消费异步流

异步流为以异步方式检索或生成元素的数据源提供自然的编程模型。本教程将创建这样的数据源,异步消费它,理解取消和上下文捕获,并比较异步接口与传统同步序列的适用方式。

先决条件

准备可运行 .NET 的环境和 C# 编译器,例如 Visual Studio 2022 或 .NET SDK。教程假定读者熟悉 C#、.NET,以及 Visual Studio 或 .NET CLI。

示例访问 GitHub GraphQL 接口,需要 GitHub 访问令牌。原文列出的权限是 repo:status 和 public_repo。应依据 GitHub 当前令牌类型和访问目标核对所需权限,并安全保存令牌。持有令牌的软件可以以它的权限调用 GitHub API。

运行起始应用

起始项目是控制台应用,使用 GitHub GraphQL 获取 dotnet/docs 仓库的最近问题。其 Main 方法如下:

static async Task Main(string[] args)
{
    //Follow these steps to create a GitHub Access Token
    // https://help.github.com/articles/creating-a-personal-access-token-for-the-command-line/#creating-a-token
    //Select the following permissions for your GitHub Access Token:
    // - repo:status
    // - public_repo
    // Replace the 3rd parameter to the following code with your GitHub access token.
    var key = GetEnvVariable("GitHubKey",
    "You must store your GitHub key in the 'GitHubKey' environment variable",
    "");

    var client = new GitHubClient(new Octokit.ProductHeaderValue("IssueQueryDemo"))
    {
        Credentials = new Octokit.Credentials(key)
    };

    var progressReporter = new progressStatus((num) =>
    {
        Console.WriteLine($"Received {num} issues in total");
    });
    CancellationTokenSource cancellationSource = new CancellationTokenSource();

    try
    {
        var results = await RunPagedQueryAsync(client, PagedIssueQuery, "docs",
            cancellationSource.Token, progressReporter);
        foreach(var issue in results)
            Console.WriteLine(issue);
    }
    catch (OperationCanceledException)
    {
        Console.WriteLine("Work has been cancelled");
    }
}

原文允许通过 GitHubKey 环境变量传入令牌,或替换 GetEnvVariable 调用的最后一个参数。共享代码时应使用安全的运行配置,不能把令牌写入源代码或提交到公共、共享仓库。

创建 GitHub 客户端后,Main 创建进度报告对象和取消令牌,随后调用 RunPagedQueryAsync,取回最多 250 个最近问题,等任务完成后再输出。

原文描述的运行行为是:每取回一页,会报告一次进度;两页请求之间可能看到停顿;只有全部十页取回后,问题才集中显示。这是教程预期行为,本稿未执行该应用。

检查原有实现

RunPagedQueryAsync 的实现如下:

private static async Task<JArray> RunPagedQueryAsync(GitHubClient client, string queryText, string repoName, CancellationToken cancel, IProgress<int> progress)
{
    var issueAndPRQuery = new GraphQLRequest
    {
        Query = queryText
    };
    issueAndPRQuery.Variables["repo_name"] = repoName;

    JArray finalResults = new JArray();
    bool hasMorePages = true;
    int pagesReturned = 0;
    int issuesReturned = 0;

    // Stop with 10 pages, because these are large repos:
    while (hasMorePages && (pagesReturned++ < 10))
    {
        var postBody = issueAndPRQuery.ToJsonText();
        var response = await client.Connection.Post<string>(new Uri("https://api.github.com/graphql"),
            postBody, "application/json", "application/json");

        JObject results = JObject.Parse(response.HttpResponse.Body.ToString()!);

        int totalCount = (int)issues(results)["totalCount"]!;
        hasMorePages = (bool)pageInfo(results)["hasPreviousPage"]!;
        issueAndPRQuery.Variables["start_cursor"] = pageInfo(results)["startCursor"]!.ToString();
        issuesReturned += issues(results)["nodes"]!.Count();
        finalResults.Merge(issues(results)["nodes"]!);
        progress?.Report(issuesReturned);
        cancel.ThrowIfCancellationRequested();
    }
    return finalResults;

    JObject issues(JObject result) => (JObject)result["data"]!["repository"]!["issues"]!;
    JObject pageInfo(JObject result) => (JObject)issues(result)["pageInfo"]!;
}

首先,方法使用 GraphQLRequest 构造 POST 请求:

public class GraphQLRequest
{
    [JsonProperty("query")]
    public string? Query { get; set; }

    [JsonProperty("variables")]
    public IDictionary<string, object> Variables { get; } = new Dictionary<string, object>();

    public string ToJsonText() =>
        JsonConvert.SerializeObject(this);
}

ToJsonText 将请求对象序列化为 JSON 字符串。JSON 序列化处理格式与必要的转义,供请求正文使用。

分页算法从最新问题向更早的问题遍历。每页请求 25 项,依据 pageInfo 中的 hasPreviousPage 和 startCursor 决定是否继续,以及上一页的起点。原文说明文字有 hasPreviousPages 的复数写法,代码和实际字段使用单数 hasPreviousPage。

问题位于 nodes 数组。方法把每页节点合并到 finalResults,然后报告累计数量并检查取消状态。如果已请求取消,ThrowIfCancellationRequested 会抛出 OperationCanceledException。

这里最大的限制是必须为所有返回问题分配存储,消费方只能等全部结果到齐。示例因此限制在 250 项。此外,进度与取消涉及 CancellationTokenSource、CancellationToken 和回调,需要追踪请求取消和响应取消的位置,首次阅读更复杂。

异步流的接口

使用 async 声明异步迭代器,在其中用 yield return 逐项返回结果;调用方用 await foreach 消费,就像同步序列使用 foreach 一样。

语言支持依赖三个在 .NET Standard 2.1 中加入、在 .NET Core 3.0 中实现的接口:

它们分别对应同步的 IEnumerable<T>、IEnumerator<T> 与 IDisposable。其中还使用 ValueTask,它的 API 与 Task 类似;原文说明这些接口出于性能考虑采用它。

改写为异步流

将 RunPagedQueryAsync 的返回类型改为 IAsyncEnumerable<JToken>,先移除取消和进度参数:

private static async IAsyncEnumerable<JToken> RunPagedQueryAsync(GitHubClient client,
    string queryText, string repoName)

原先每页取回后执行三行代码:

finalResults.Merge(issues(results)["nodes"]!);
progress?.Report(issuesReturned);
cancel.ThrowIfCancellationRequested();

替换为逐项产出:

foreach (JObject issue in issues(results)["nodes"]!)
    yield return issue;

同时删除 finalResults 变量及循环后的返回该集合的 return 语句。完整异步迭代器如下:

private static async IAsyncEnumerable<JToken> RunPagedQueryAsync(GitHubClient client,
    string queryText, string repoName)
{
    var issueAndPRQuery = new GraphQLRequest
    {
        Query = queryText
    };
    issueAndPRQuery.Variables["repo_name"] = repoName;

    bool hasMorePages = true;
    int pagesReturned = 0;
    int issuesReturned = 0;

    // Stop with 10 pages, because these are large repos:
    while (hasMorePages && (pagesReturned++ < 10))
    {
        var postBody = issueAndPRQuery.ToJsonText();
        var response = await client.Connection.Post<string>(new Uri("https://api.github.com/graphql"),
            postBody, "application/json", "application/json");

        JObject results = JObject.Parse(response.HttpResponse.Body.ToString()!);

        int totalCount = (int)issues(results)["totalCount"]!;
        hasMorePages = (bool)pageInfo(results)["hasPreviousPage"]!;
        issueAndPRQuery.Variables["start_cursor"] = pageInfo(results)["startCursor"]!.ToString();
        issuesReturned += issues(results)["nodes"]!.Count();

        foreach (JObject issue in issues(results)["nodes"]!)
            yield return issue;
    }

    JObject issues(JObject result) => (JObject)result["data"]!["repository"]!["issues"]!;
    JObject pageInfo(JObject result) => (JObject)issues(result)["pageInfo"]!;
}

接着,找到 Main 中原来的集合消费代码:

var progressReporter = new progressStatus((num) =>
{
    Console.WriteLine($"Received {num} issues in total");
});
CancellationTokenSource cancellationSource = new CancellationTokenSource();

try
{
    var results = await RunPagedQueryAsync(client, PagedIssueQuery, "docs",
        cancellationSource.Token, progressReporter);
    foreach(var issue in results)
        Console.WriteLine(issue);
}
catch (OperationCanceledException)
{
    Console.WriteLine("Work has been cancelled");
}

用 await foreach 替换:

int num = 0;
await foreach (var issue in RunPagedQueryAsync(client, PagedIssueQuery, "docs"))
{
    Console.WriteLine(issue);
    Console.WriteLine($"Received {++num} issues in total");
}

IAsyncEnumerator<T> 实现 IAsyncDisposable,因此结束枚举时可以异步释放迭代器。原文用下面的等价思路解释循环:

int num = 0;
var enumerator = RunPagedQueryAsync(client, PagedIssueQuery, "docs").GetAsyncEnumerator();
try
{
    while (await enumerator.MoveNextAsync())
    {
        var issue = enumerator.Current;
        Console.WriteLine(issue);
        Console.WriteLine($"Received {++num} issues in total");
    }
} finally
{
    if (enumerator != null)
        await enumerator.DisposeAsync();
}

默认情况下,流元素在捕获的上下文中处理。如果希望禁止捕获,可以使用 TaskAsyncEnumerableExtensions.ConfigureAwait。同步上下文与捕获机制详见消费基于任务的异步模式。

取消异步流

可以在异步迭代器签名中添加具有 [EnumeratorCancellation] 属性的取消令牌:

private static async IAsyncEnumerable<JToken> RunPagedQueryAsync(GitHubClient client,
    string queryText, string repoName, [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    var issueAndPRQuery = new GraphQLRequest
    {
        Query = queryText
    };
    issueAndPRQuery.Variables["repo_name"] = repoName;

    bool hasMorePages = true;
    int pagesReturned = 0;
    int issuesReturned = 0;

    // Stop with 10 pages, because these are large repos:
    while (hasMorePages && (pagesReturned++ < 10))
    {
        var postBody = issueAndPRQuery.ToJsonText();
        var response = await client.Connection.Post<string>(new Uri("https://api.github.com/graphql"),
            postBody, "application/json", "application/json");

        JObject results = JObject.Parse(response.HttpResponse.Body.ToString()!);

        int totalCount = (int)issues(results)["totalCount"]!;
        hasMorePages = (bool)pageInfo(results)["hasPreviousPage"]!;
        issueAndPRQuery.Variables["start_cursor"] = pageInfo(results)["startCursor"]!.ToString();
        issuesReturned += issues(results)["nodes"]!.Count();

        foreach (JObject issue in issues(results)["nodes"]!)
            yield return issue;
    }

    JObject issues(JObject result) => (JObject)result["data"]!["repository"]!["issues"]!;
    JObject pageInfo(JObject result) => (JObject)issues(result)["pageInfo"]!;
}

EnumeratorCancellationAttribute使编译器生成的迭代器代码能够把传给 GetAsyncEnumerator 的令牌传入方法体。实现方随后可以检查令牌,响应取消请求。

调用方通过 WithCancellation传递令牌:

private static async Task EnumerateWithCancellation(GitHubClient client)
{
    int num = 0;
    var cancellation = new CancellationTokenSource();
    await foreach (var issue in RunPagedQueryAsync(client, PagedIssueQuery, "docs")
        .WithCancellation(cancellation.Token))
    {
        Console.WriteLine(issue);
        Console.WriteLine($"Received {++num} issues in total");
    }
}

上面的源代码完整保留了原文的示例,但添加参数并不等于完成取消实现:这个迭代器方法体没有检查 cancellationToken,也没有把它传给底层 Post 请求;消费示例创建令牌源后也未展示触发取消。因此,不能据此声称正在进行的网络请求会被中断。实际应用需在生产者和请求层落实取消,并按需处理取消异常。

完成后的项目也位于 dotnet/docs 仓库。

比较完成后的应用

原文建议重新运行应用,与起始项目比较。异步流可以在第一页结果可用时立即枚举,之后每次分页请求仍可能有等待,但不必为显示第一页先保存所有页面。每项输出与累计数量可以直接放在 await foreach 中,省去独立进度回调。

调用方停止枚举,可以结束继续消费;这与令牌取消底层请求是两个层次。如果实际取消引发异常,仍应按应用需求处理,不能把原文“不需要 try/catch”理解为所有取消情形都无需异常处理。

从代码结构可以看出,生产者不再维护完整结果集合,消费方决定是否另行保存。这是结构上的存储需求变化,本文未实测内存或性能。

运行起始与完成项目可自行观察差异。练习结束后,原文建议删除为本教程创建的 GitHub 访问令牌,防止它被继续滥用。

本例从分页网络 API 逐项读取。异步流也适用于行情或传感器等持续产生数据的源:MoveNextAsync 等待下一项可用,再交给消费方。


来源:Microsoft 文档团队,原文,页面标示更新于 2025-06-19,核对日期 2026-10-03。文档采用 CC BY 4.0,仓库许可。本版整理中文措辞,指出字段拼写与取消示例的实现边界,代码逐块保留原样,未运行。代码采用 MIT,完整通知如下。

The MIT License (MIT) Copyright (c) Microsoft Corporation

Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the “Software”), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions: The above copyright notice and this permission notice shall be included in all copies or substantial portions of the Software. THE SOFTWARE IS PROVIDED “AS IS”, WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容