将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 尚未被支持。 如果将流分析指向Azure Cosmos DB使用其他 API 创建的帐户,则可能无法正确存储数据。 当你使用Azure Cosmos DB作为输出时,将作业设置为兼容性等级1.2。

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

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

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

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

使用流分析进行插入或更新操作

通过将 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 的文档),请重命名查询中的 ID 字段(使用“AS”关键字)。 让我们Azure Cosmos DB创建 ID 字段或将 ID 替换为另一列的值(通过使用 AS 关键字或使用 Document ID 设置)。

Azure Cosmos DB中的数据分区

Azure Cosmos DB根据工作负荷自动缩放分区。 使用 无限 容器来分区数据。 当 Stream Analytics 写入无限制的容器时,它使用的并行写入器数量与先前查询步骤或输入分区方案的数量一致。

注意

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,upsert 行为会改变,插入或替换文档。

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

以下示例演示从同一 Azure 事件中心 输入读取的两个相同的 Stream Analytics 作业。 两个流分析作业均完全分区,通过直接查询并写入到相同的多个 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 配合使用时如何横向扩展以接收每秒 1,000 到 10,000 条消息,请尝试此 Azure 示例项目

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

用于 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 的容器。
文档 ID 可选。 输出事件中的列名,作为插入或更新操作的唯一键。 如果将其留空,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 中的唯一密钥约束。