对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

Atlas Stream Processing操作

本指南介绍了管理Atlas Stream Processing 的最佳实践,包括基本工作流程的特定程序。

每个流处理器都依赖于与数据源和接收器的连接。 Atlas Stream Processing使用 TLS/SSL 对传输中的所有数据进行加密,但基本连接仍会通过公共互联网传递数据。为了使用Atlas Stream Processing实现最佳性能和安全性,请考虑本节中描述的实践。

为确保Atlas Stream Processing与私有外部系统(例如自管理Apache Kafka集群或私有 API)之间的安全通信,请使用VPC或 VNet 对等互连。这些连接可防止您的数据暴露在公共互联网上。请注意,除了建立对等连接之外,您还必须为应用程序的VPC显式配置路由表,以将流量定向到Atlas VPC CIDR区块。要学习;了解有关如何配置路由表的更多信息,请参阅外部提供商的VPC文档。

为了在同一组织内的不同Atlas项目之间路由流媒体数据, Atlas Stream Processing支持跨项目连接。跨项目连接传输数据时不会将集群暴露于私有互联网,从而确保隐私,无需手动配置VPC对等互连。

Atlas Stream Processing支持以下 AWS 集成:

  • S3

  • Kinesis Data Streams

  • Lambda (如 $externalFunction

其中每个集成都与统一 AWS 访问兼容。配置 统一访问 无需在MongoDB中存储静态访问权限ID 或密钥凭证,并且是所有 AWS 集成推荐的身份验证模型。

Atlas Stream Processing支持使用 OIDC 和JSON Web Tokens 对Apache Kafka代理进行身份验证。但是,此功能目前仅支持可通过公共网络访问的身份提供程序。客户VPC中的身份提供商不支持此功能。

Atlas Stream Processing提供了各种机制来确保在常见故障情况下提供可靠、稳健的服务,包括检查点、死信队列和故障转移处理器。

Atlas Stream Processing使用 30 秒心跳监控流处理器的运行状况。如果流处理器未能发送心跳超过 30 秒, Atlas Stream Processing会自动重启该处理器。

Atlas Stream Processing自动管理 AWS Kinesis分片的集成的拓扑结构更改。

您可以配置Atlas Stream Processing ,将未处理的数据路由到死信队列(DLQ)(DLQ)。然后,您可以在 DLQ 中查看未处理的记录,以修复处理故障。

检查 DLQ 时,请查找以下常见处理错误:

  • 有效负载解析

    传入记录格式不正确,无法解析。这通常是由于无效的JSON造成的。

  • 表达式求值

    处理管道中的表达式无法正确计算。

    例子

    您的管道包含用于删除Kafka记录的 $emit.config.tombstoneWhen 配置字段,但给定文档缺少要评估的字段。

  • 大小限制

    传出记录超过目标接收器的大小限制,例如 AWS Kinesis的 1 MB

  • 密钥生成

    处理器在为 Sink 生成分区键时遇到错误。

  • 延迟数据

    您的处理器使用 Windows,并且给定文档在配置的允许延迟期到期后到达。

分析处理失败的原因后,修改上游数据源,在数据进入Atlas Stream Processing管道之前更正数据,或者修改管道以处理边缘情况。

在使用 DLQ 中的记录来指南对Atlas Stream Processing配置的更正后,请考虑对这些记录使用以下修复选项:

  • 丢弃

    如果某条记录严重错误或不相关(例如时间高度敏感的系统中的数据非常晚的情况),则将其从集合中删除。

  • 手动更正

    对于存在轻微JSON错误或关键任务的记录,请手动调整字段并将其直接插入目标接收器。

  • 自动重新处理:

    对于大量可恢复的相关数据,请创建从 DLQ 中读取的从节点(secondary node from replica set)脚本或额外的流处理器,并在将清理后的数据路由到目标接收器之前根据已知故障条件进行转换。

Atlas Stream Processing提供故障转移处理器,防止在全区域系统出现故障事件服务中断。

以下原则定义了强大的故障转移配置,以保护您的Atlas Stream Processing工作负载。

  • 在多个区域部署Atlas集群。

    为确保流处理器在区域故障转移事件中继续运行,流处理器连接的任何Atlas 集群都必须是多区域集群。这一原则可确保您的源连接或汇连接在主节点 (primary node in the replica set)地区关闭时保持活动状态。

  • 为每个地区创建Atlas Stream Processing工作区。

    部署多区域集群后,为每个地区创建一个工作区。这一原则可最大限度地减少流处理器与其支持的Atlas数据库之间的连接延迟。

  • 配置特定于连接的灾难恢复

    为了确保流处理器的服务连续性,它所依赖的任何外部提供程序也必须启用故障转移。确保每个外部提供商源或接收器都根据其文档进行正确配置。

    例子

    如果您的流处理器使用Kafka集群作为源或接收器,请将集群配置为保留消费者偏移量,以启用Atlas Stream Processing能够从消费者群组偏移量恢复处理。

  • 执行例行演练。

    通过定期执行故障转移转移试运行来验证故障转移配置。这些试提供服务既可确保您的配置正确连接到网络和加载检查点,又能训练故障转移工作流程,以便在需要时更轻松地执行。考虑建立定期的故障转移转移试运行节奏,并在每次更改故障转移配置时执行试运行。

确保最佳实践就位后,创建流处理器并验证故障转移配置。

Atlas Stream Processing仅支持具有以下配置的处理器的自动区域故障转移:

有关自动故障转移的更多信息,请参阅故障转移处理器。

但是,您可以配置Atlas Stream Processing架构,支持对任意源连接和接收器连接组合进行手动故障转移。

您可以在工作区或流处理器级别启动强制故障转移。

要启动流处理工作区故障转移,请按照触发流处理工作区故障转移中描述的步骤进行操作。

要启动单个流处理器故障转移,请按照为一个流处理器启动故障转移中所述的步骤进行操作。

如果发生区域服务中断事件:

2

按照为一个流处理器启动故障转移中描述的步骤进行操作。为确保服务从上次使用的流事件恢复,请将您之前记下的时间戳作为 startAtOperationTime 参数传递。

使用Apache Kafka源的处理器根据消费者群组偏移量进行恢复。

3

当主节点 (primary node in the replica set)地区重新在线时,主节点 (primary node in the replica set)处理器将自动重启,从而导致重复处理。创建并运行脚本以持续检查主节点 (primary node in the replica set)地区的状态,并在主区域恢复在线时停止该地区中的所有处理器。

4

完成故障转移后,请定期检查工作区主节点 (primary node in the replica set)地区的状态。您可以在MongoDB Cloud 状态页面或云提供商的状态页面上查看此信息。您还可以配置Atlas Stream Processing警报。如果使用自动主节点 (primary node in the replica set)处理器抑制,则可以将脚本配置为在检测到主节点 (primary node in the replica set)处理器自动重新启动时返回通知。

5

停止故障转移流处理器。请注意每次停止的时间戳。

6

重新启动主节点 (primary node in the replica set)地区中的流处理器。为确保服务从上次使用的流事件恢复,请将您之前记下的时间戳作为 startAtOperationTime 参数传递。

使用Apache Kafka源的处理器根据消费者群组偏移量进行恢复。

在为区域故障转移或触发手动故障转移系统配置系统时,请考虑以下事项:

Atlas Stream Processing提供各种可观测性工具来评估工作区、连接和处理器的性能和运行状况。此外,要追踪区域服务状态和中断,您可以访问MongoDB Cloud 状态页面。

Atlas Stream Processing为每个流处理器配置资源。要学习;了解更多信息,请参阅《层级选择指南》。

对于在Atlas集群上运行的Atlas Stream Processing处理器,请确保根据管道I/O 需求按比例扩展处理器和集群。特别是,评估执行大量 $lookup$merge 操作的处理器,将其作为集群升级的候选处理器。

Atlas Stream Processing可以使用Atlas 集群$source 来消费变更流事件。为了确保此类流处理器的最佳性能和可靠性,查看以下注意事项:

  • Oplog 窗口大小

    流处理器在停止后利用 $source 集群的oplog来恢复服务。如果处理器关闭的时间长于oplog保留期,则恢复令牌可能会点截断的数据,从而导致 ChangeStreamHistoryLost 错误。

    通过使oplog大小与流处理系统的可接受停机时间保持一致来降低此风险。

  • 恢复选项

    Atlas Stream Processing支持恢复变更流 $source``s from either the last stored checkpoint with ``resumeFromCheckpoint 或使用 startAtOperationTime 恢复特定时间。为了避免处理的数据出现空白或重复,请围绕在给定场景中使用的恢复选项定义标准化策略。

  • 待办事项摄取

    在大型命名空间上运行的流处理器可能会经历很长的初始追赶处理期。请考虑缩小流处理器命名空间的范围或选择较新的点,以缩短初始追赶期。

  • 变更流图像

    审核、补偿事务和撤消操作需要访问权限预处理文档状态。当您使用Atlas作为源时,启用文档前像和后像以支持这些操作。当您推出对文档前图像或后图像的支持时,请验证整个系统的覆盖范围和稳定性。

  • ``$replaceRoot`` 阶段

    变更流源返回将操作元数据与受影响文档内容捆绑在一起的文档。将这些文档写入接收器需要使用者使用变更事件模式,而不是源文档的业务逻辑。除非您需要元数据,否则请在进一步处理之前使用 $replaceRoot 阶段提升源内容子文档。

  • 面向连接的管道

    Atlas Stream Processing支持 $lookup 阶段,以启用流媒体数据的类似联接的丰富。这些操作会增加处理延迟。考虑将Atlas Stream Processing管道限制为仅执行在写入接收器之前必须执行的 $lookup 操作,为预计算视图保留更复杂的扩展。

  • 过时检查点

    检查点仅包含事件历史记录,不包含处理逻辑。当您从检查点恢复操作时,流处理器会将当前管道的逻辑应用于旧事件,这可能会导致现有下游状态与新输出之间不匹配。如果您需要显着不同的处理语义,请考虑引入新的流处理器,而不是从检查点恢复。

  • 重新启动规划

    重启期间,流处理器的新事件吞吐量可能会降低。为了提高系统弹性,安排会重新启动以考虑暂时减少的容量。

  • 故障转移规划

    制定跨区域共享流处理器检查点、选择安全恢复点和验证输出一致性的标准程序,以确保区域故障转移事件不会导致数据质量问题。