你当前正在访问 Microsoft Azure Global Edition 技术文档网站。 如果需要访问由世纪互联运营的 Microsoft Azure 中国技术文档网站,请访问 https://docs.azure.cn

Azure 流分析输出到 Azure Cosmos DB

Azure Cosmos DB 在 Azure 流分析 中的输出会将流处理结果写成 JSON 文档到 Azure Cosmos DB 容器中。 它支持对非结构化 JSON 数据进行数据归档和低延迟查询。 了解该输出的行为方式有助于你根据场景所需的吞吐量、一致性和分区对其进行配置。

作为输出目标的 Azure Cosmos DB 的基础知识

Stream Analytics 中的 Azure Cosmos DB 输出会将你的流处理结果写成 JSON 输出到你的 Azure Cosmos DB 容器中。 如果不熟悉 Azure Cosmos DB,请参阅 Azure Cosmos DB 文档了解入门知识。

Stream Analytics 仅通过 SQL API 连接 Azure Cosmos DB。 其他 Azure Cosmos DB API 尚未被支持。 如果使用其他 API 将流分析指向已创建的 Azure Cosmos DB 帐户,则可能无法正确存储数据。 当你使用Azure Cosmos DB作为输出时,将作业设置为兼容性等级1.2。

流分析不会在数据库中创建容器。 相反,你需要提前创建它们。 然后,你可以控制 Azure Cosmos DB 容器的计费成本。 还可以使用 Azure Cosmos DB APIs 直接调整容器的性能、一致性和容量。 以下部分详细介绍了 Azure Cosmos DB 的一些容器选项。

调整一致性、可用性和延迟

为了满足你的应用需求,在 Azure Cosmos DB 中微调数据库和容器,并在一致性、可用性、延迟和吞吐量之间做出权衡。

根据你的场景对读取一致性和读写延迟之间的要求,为数据库帐户选择一种一致性级别。 为了提高吞吐量,可以在容器上扩展请求单元(RU)。 另外,默认情况下,Azure Cosmos DB 会对容器的每个 CRUD 操作启用同步索引。 这个选项是控制 Azure Cosmos DB 读写性能的另一种有用方式。 有关详细信息,请参阅更改数据库和查询的一致性级别一文。

从流分析进行 Upsert 操作

通过将 Stream Analytics 与 Azure Cosmos DB 集成,你可以根据给定的 文档 ID 列在容器中插入或更新记录。 此操作也称为“更新插入”。 流分析使用乐观 Upsert 方法。 即仅当由于文档 ID 冲突而插入失败时才进行更新。

通过使用兼容性级别 1.0,Stream Analytics 将此更新作为 PATCH 操作执行,因此支持文档的部分更新。 流分析添加新属性或增量替换现有属性。 但是,JSON 文档中数组属性值的更改会导致覆盖整个数组。 也就是说,不会合并数组。

使用兼容性级别 1.2 后,upsert 行为将变为插入或替换文档。 关于兼容性级别 1.2 的部分,后面将进一步描述这种行为。

如果收到的 JSON 文档已有 ID 字段,Azure Cosmos DB 会自动将该字段作为文档 ID 列使用。 Stream Analytics 会据此处理任何后续写入,从而导致以下情况之一:

  • 唯一标识符导致插入。
  • 重复的 ID 和设置为“ID”的“文档 ID”导致更新插入。
  • 在第一个文档之后,重复的 ID 和未设置的“文档 ID”导致错误。

如果要保存“所有”文档(包括具有重复 ID 的文档),请重命名查询中的 ID 字段(使用“AS”关键字)。 让 Azure Cosmos DB 创建 ID 字段,或用另一列的值替换 ID(使用“AS”关键字,或者使用“文档 ID”设置)。

Azure Cosmos DB 中的数据分区

Azure Cosmos DB 会根据工作负载自动缩放分区。 使用 无限 容器来分区数据。 当流分析写入多个无限制的容器时,它会根据先前查询步骤或输入分区方案使用相应数量的并行写入器。

注意

Azure 流分析仅支持顶级分区键的无限制容器。 例如,支持 /region。 嵌套分区键(例如) /region/name不被支持。

你可能会收到以下警告,具体取决于你选择的分区键:

CosmosDB Output contains multiple rows and just one row per partition key. If the output latency is higher than expected, consider choosing a partition key that contains at least several hundred records per partition key.

选择一个具有多个不同值的分区键属性,并且能均匀分配工作负载到这些值。 作为分区的一个自然现象,单个分区的最大吞吐量限制了涉及相同分区键的请求。

属于同一分区键值的文档的存储大小上限为 20 GB(物理分区大小上限为 50 GB)。 理想的分区键应当作为查询中的过滤器频繁出现,并且具有足够的基数以确保你的解决方案具有可扩展性。

用于流分析查询和 Azure Cosmos DB 的分区键不需要完全相同。 对于完全并行拓扑,可以使用 Input Partition keyPartitionId作为 Stream Analytics 查询的分区键,但这个选项可能不是 Azure Cosmos DB 容器的分区键的推荐选择。

分区键也是 Azure Cosmos DB 中存储过程和触发器的事务边界。 选择分区键,使事务中出现的文档共享相同的分区键值。 如需详细了解如何选择分区键,请参阅 Azure Cosmos DB 中的分区一文。

对于固定容量的 Azure Cosmos DB 容器,Stream Analytics 在其容量用尽后无法对其进行纵向扩展或横向扩展。 这些集合的大小上限为 10 GB,吞吐量上限为 10,000 RU/秒。 若要将数据从固定的容器迁移到无限制容器(例如,吞吐量至少为 1,000 RU/秒,且具有分区键),请使用数据迁移工具更改源库

将逐步弃用写入多个固定容器的功能。 不要用它来扩展你的流分析工作。

使用兼容性级别 1.2 改进了吞吐量

通过使用1.2级兼容性,Stream Analytics 支持原生集成批量写入 Azure Cosmos DB。 通过这种集成,Stream Analytics 能够有效地写入 Azure Cosmos DB,同时最大化吞吐量并高效处理限速请求。

新兼容性级别改变了更新插入行为,提供一种改进的写入机制。 使用 1.2 之前的级别时,upsert 操作会插入或合并文档。 使用 1.2 时,更新插入行为将更改为插入或替换文档。

使用 1.2 之前的级别时,流分析使用自定义存储过程,按分区键将文档批量更新插入到 Azure Cosmos DB。 在该过程中,Stream Analytics 会将一个批次作为一笔事务写入。 即使单条记录出现瞬态错误(限速),Stream Analytics 也必须重试整个批次。 这种行为让即使是合理的限速场景也会变慢。

以下示例显示了从同一 Azure 事件中心输入读取数据的两个相同流分析作业。 这两个流分析作业已使用直通查询完全分区,并写入到相同的 Azure Cosmos DB 容器。 左侧的指标来自配置了兼容性级别 1.0 的作业。 右侧的指标来自配置为使用 1.2 的作业。 Azure Cosmos DB 容器的分区键是一个来自输入事件的 GUID,它是唯一的。

显示流分析指标比较的屏幕截图。

Event Hubs 的传入事件速率是 Azure Cosmos DB 容器(20,000 RU)配置可接收速率的两倍,因此可以预期 Azure Cosmos DB 会出现限流。 但是,使用版本 1.2 的作业一贯以更高的吞吐量(输出事件数/分钟)写入,并且其平均 SU% 利用率更低。 在你的环境中,这种差异还取决于其他几个因素。 这些因素包括:事件格式的选择、输入事件/消息大小、分区键和查询。

显示 Azure Cosmos DB 指标比较的屏幕截图。

通过使用 1.2 版本,Stream Analytics 更智能地利用了 Azure Cosmos DB 中 100% 的可用吞吐量,几乎没有因限速或速率限制而产生的重新提交。 对于其他工作负荷(例如,同时在容器上运行的查询),此行为可以提供更好的体验。 如需了解 Azure Cosmos DB 作为接收器(每秒接收 1000 到 10000 条消息)如何横向扩展流分析,请尝试此 Azure 示例项目

使用 1.0 和 1.1 版本时,Azure Cosmos DB 输出的吞吐量是相同的。 强烈建议对 Azure Cosmos DB 的流分析使用兼容性级别 1.2。

JSON 输出的 Azure Cosmos DB 设置

当你在Stream Analytics中将Azure Cosmos DB配置为输出时,以下属性定义了输出。

显示 Azure Cosmos DB 输出流的信息字段的屏幕截图。

字段 说明
输出别名 用于在流分析查询中引用此输出的别名。
订阅 Azure 订阅。
帐户 ID Azure Cosmos DB 帐户的名称或终结点 URI。
帐户密钥 Azure Cosmos DB 帐户的共享访问密钥。
数据库 Azure Cosmos DB 数据库名称。
容器名称 容器名称,如 MyContainer。 必须存在名为 MyContainer 的容器。
文档编号 可选。 输出事件中用作插入或更新操作唯一键的列名。 如果将其留空,Stream Analytics 则会插入所有事件,且不提供更新选项。

配置 Azure Cosmos DB 输出后,可以在查询中将其用作 INTO 语句的目标。 当你以这种方式使用 Azure Cosmos DB 输出时,必须显式设置分区键

输出记录必须包含一个区分大小写的列,该列以 Azure Cosmos DB 中的分区键命名。 若要实现更大的并行化,该语句可能需要使用同一列的 PARTITION BY 子句

下面是一个示例查询:

    SELECT TollBoothId, PartitionId
    INTO CosmosDBOutput
    FROM Input1 PARTITION BY PartitionId

错误处理和重试

如果流分析将事件发送到 Azure Cosmos DB 时出现了暂时性故障、服务不可用或受到限制,流分析将无限期重试以成功完成操作。 但它不会针对未经授权(HTTP 错误代码 401)、未找到(HTTP 错误代码 404)、禁止访问(HTTP 错误代码 403)或错误请求(HTTP 错误代码 400)失败进行重试。

导致 Azure Cosmos DB 输出失败的常见问题

多种情况可能导致 Azure Cosmos DB 输出失败。 Stream Analytics 的输出数据可能违反了容器的唯一索引约束,或者 PartitionKey 该列可能不存在,或者 Id 该列可能不存在。 有关唯一索引约束的更多信息,请参见 Azure Cosmos DB 中的唯一密钥约束。