本指南介绍了管理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项目之间路由流媒体数据, Atlas Stream Processing支持跨项目连接。跨项目连接传输数据时不会将集群暴露于私有互联网,从而确保隐私,无需手动配置VPC对等互连。
无密码 AWS 身份验证
Atlas Stream Processing支持以下 AWS 集成:
S3
Kinesis Data Streams
Lambda (如
$externalFunction)
其中每个集成都与统一 AWS 访问兼容。配置 统一访问 无需在MongoDB中存储静态访问权限ID 或密钥凭证,并且是所有 AWS 集成推荐的身份验证模型。
Kafka OIDC
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 生成分区键时遇到错误。
分析处理失败的原因后,修改上游数据源,在数据进入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 源和汇聚器
Atlas source and Apache Kafka sink
但是,您可以配置Atlas Stream Processing架构,支持对任意源连接和接收器连接组合进行手动故障转移。
启动强制故障转移
您可以在工作区或流处理器级别启动强制故障转移。
响应故障转移事件
如果发生区域服务中断事件:
验证服务中断情况。
当您发现服务中断的症状时,查看您首选的状态报告以验证该事件。您可以在MongoDB Cloud 状态页面或云提供商的状态页面上查看此信息。您可以配置Atlas Stream Processing警报。记下失败事件的时间戳。
启动故障转移处理器。
按照为一个流处理器启动故障转移中描述的步骤进行操作。为确保服务从上次使用的流事件恢复,请将您之前记下的时间戳作为 startAtOperationTime 参数传递。
使用Apache Kafka源的处理器根据消费者群组偏移量进行恢复。
验证服务是否已恢复。
完成故障转移后,请定期检查工作区主节点 (primary node in the replica set)地区的状态。您可以在MongoDB Cloud 状态页面或云提供商的状态页面上查看此信息。您还可以配置Atlas Stream Processing警报。如果使用自动主节点 (primary node in the replica set)处理器抑制,则可以将脚本配置为在检测到主节点 (primary node in the replica set)处理器自动重新启动时返回通知。
重新启动主节点 (primary node in the replica set)处理器。
重新启动主节点 (primary node in the replica set)地区中的流处理器。为确保服务从上次使用的流事件恢复,请将您之前记下的时间戳作为 startAtOperationTime 参数传递。
使用Apache Kafka源的处理器根据消费者群组偏移量进行恢复。
故障转移行为
在为区域故障转移或触发手动故障转移系统配置系统时,请考虑以下事项:
发生区域故障转移时, Atlas Stream Processing会将故障转移地区中的不活动故障转移处理器提升为活动处理器。这些处理器从最后一个可用检查点恢复流处理。
Atlas Stream Processing检查点的出现时间保证不超过 10 分钟。但是,故障转移处理器会重新处理主节点 (primary node in the replica set)处理器自上一个启动检查点以来处理过的所有数据。故障转移ACID 一致性保证至少一次处理语义。
如果故障转移处理器在区域故障转移事件期间无法启动,它仍会承担活动处理器的角色。然后,您可以排除根本的原因,并在故障转移地区照常重新启动处理器。
可观察性
Atlas Stream Processing提供各种可观测性工具来评估工作区、连接和处理器的性能和运行状况。此外,要追踪区域服务状态和中断,您可以访问MongoDB Cloud 状态页面。
扩展Atlas Stream Processing
Atlas Stream Processing为每个流处理器配置资源。要学习;了解更多信息,请参阅《层级选择指南》。
对于在Atlas集群上运行的Atlas Stream Processing处理器,请确保根据管道I/O 需求按比例扩展处理器和集群。特别是,评估执行大量 $lookup 或 $merge 操作的处理器,将其作为集群升级的候选处理器。
变更流注意事项
Atlas Stream Processing可以使用Atlas 集群$source 来消费变更流事件。为了确保此类流处理器的最佳性能和可靠性,查看以下注意事项:
恢复选项
Atlas Stream Processing支持恢复变更流
$source``s from either the last stored checkpoint with ``resumeFromCheckpoint或使用startAtOperationTime恢复特定时间。为了避免处理的数据出现空白或重复,请围绕在给定场景中使用的恢复选项定义标准化策略。待办事项摄取
在大型命名空间上运行的流处理器可能会经历很长的初始追赶处理期。请考虑缩小流处理器命名空间的范围或选择较新的点,以缩短初始追赶期。
``$replaceRoot`` 阶段
变更流源返回将操作元数据与受影响文档内容捆绑在一起的文档。将这些文档写入接收器需要使用者使用变更事件模式,而不是源文档的业务逻辑。除非您需要元数据,否则请在进一步处理之前使用
$replaceRoot阶段提升源内容子文档。面向连接的管道
Atlas Stream Processing支持
$lookup阶段,以启用流媒体数据的类似联接的丰富。这些操作会增加处理延迟。考虑将Atlas Stream Processing管道限制为仅执行在写入接收器之前必须执行的$lookup操作,为预计算视图保留更复杂的扩展。过时检查点
检查点仅包含事件历史记录,不包含处理逻辑。当您从检查点恢复操作时,流处理器会将当前管道的逻辑应用于旧事件,这可能会导致现有下游状态与新输出之间不匹配。如果您需要显着不同的处理语义,请考虑引入新的流处理器,而不是从检查点恢复。
重新启动规划
重启期间,流处理器的新事件吞吐量可能会降低。为了提高系统弹性,安排会重新启动以考虑暂时减少的容量。
故障转移规划
制定跨区域共享流处理器检查点、选择安全恢复点和验证输出一致性的标准程序,以确保区域故障转移事件不会导致数据质量问题。