View a markdown version of this page

对亚马逊 MSK 复制器进行故障排除 - Amazon Managed Streaming for Apache Kafka

本文属于机器翻译版本。若本译文内容与英语原文存在差异,则一律以英文原文为准。

对亚马逊 MSK 复制器进行故障排除

以下信息可以帮助您解决 MSK Replicator 的问题。排查 Amazon MSK 集群的问题有关亚马逊 MSK 的其他功能,请参阅。您也可以将问题发布到 AWS re:Post。

复制器状态从 “正在创建” 变为 “失败”

创建 MSK Replicator 失败的常见原因:

  1. 验证您为目标集群提供的安全组是否具有允许流量流向目标集群安全组的出站规则,以及目标集群的安全组是否具有接受来自 Replicator 安全组流量的入站规则。

  2. 对于跨区域复制,请确认您的源集群已开启多 VPC 连接以进行 IAM 访问控制,并且集群策略是否已在源集群上设置。

  3. 验证在创建期间提供的 IAM 角色是否具有读取和写入源集群和目标集群所需的权限,包括写入主题的权限。

  4. 确认您的网络 ACL 未阻塞 MSK 复制器与集群之间的连接。

  5. 当 MSK Replicator 尝试连接时,源集群或目标群集可能尚未完全可用。这可能是由于负载过大、磁盘使用率或 CPU 使用率造成的。修复代理的问题,然后重试创建复制器。

执行上述验证后,再次创建 MSK 复制器。

复制器似乎停留在 “创建” 状态

创建 MSK 复制器最多可能需要 30 分钟。请等待 30 分钟,然后再次检查集群的状态。

Replicator 不复制数据或仅复制部分数据

  1. 使用亚马逊 CloudWatch中的AuthError指标验证您的复制器没有遇到身份验证错误。如果此指标高于 0,请检查 IAM 角色策略并确保没有为集群权限设置拒绝权限。

  2. 确认您的源集群和目标集群没有遇到问题(连接过多、磁盘满负荷或 CPU 使用率高)。

  3. 使用该KafkaClusterPingSuccessCount指标验证您的集群是否可访问。如果此指标为 0 或没有数据点,请检查网络和 IAM 角色权限。

  4. 使用该ReplicatorFailure指标验证您的复制器没有出现故障。如果大于 0,请检查 IAM 角色中的主题级权限。

  5. 验证允许列表中的正则表达式是否与要复制的主题的名称相匹配,并且主题未被拒绝列表排除。

  6. Replicator 最多可能需要 30 秒才能检测和创建新主题。如果起始位置为最新(默认),则不会复制在目标集群上创建主题之前生成的消息。

目标集群中的消息偏移量与源集群不同

MSK Replicator 使用来自源集群的消息并将其生成到目标集群,这可能会导致不同的偏移量。如果您已开启使用者组偏移量同步,MSK Replicator 将自动转换偏移量,以便在故障转移后,您的消费者可以从他们中断的地方恢复处理。

Replicator 未同步使用者组偏移量

  1. 验证数据复制是否按预期运行。

  2. 验证允许列表中的正则表达式是否与要复制的使用者组相匹配。

  3. 验证 MSK 复制器是否已在目标集群上创建了该主题。如果源集群上的消费者组仅消费了尚未复制的消息,则该使用者组不会被复制到目标集群。一旦您的消费者组开始读取新复制的消息,MSK Replicator 将自动复制该消费者组。

注意

MSK Replicator 优化了消费者群组偏移量同步,使消费者能够从主题分区接近末尾处进行读取。如果您的消费者群体在源集群上滞后,则可能会在目标集群上看到更高的延迟。当您的消费者赶上时,MSK 复制器将自动减少延迟。

复制延迟很高或持续增加

  1. 验证您的分区数量是否正确。下表显示了满足所需吞吐量的建议的最小分区数。

    吞吐量和建议的最小分区数
    吞吐量 (MB/s) 所需的最低分区
    50167
    100334
    250833
    5001666
    10003333
  2. 确认您的集群中有足够的读写容量。MSK 复制器充当源集群(出口)的使用器,也充当目标集群(入口)的生成器。除其他流量外,还要配置集群容量以支持复制流量。

  3. 复制延迟因区域对距离而异。

  4. 使用该指标确认您的复制器没有受到限制。ThrottleTime如果大于 0,则调整 Kafka 配额。请参阅使用 Kafka 配额管理吞吐量。

  5. 查看AWS 服务运行状况控制面板,了解您所在地区的 MSK 服务事件。

使用 ReplicatorFailure 指标进行故障排除

该ReplicatorFailure指标可帮助您监控和检测复制问题。非零值通常表示由消息大小限制、时间戳范围违规或记录批量大小问题导致的复制失败。如果您为复制器配置了日志传输,则可以使用传送的日志消息来识别特定的故障。有关更多详细信息,请参阅MSK 复制器日志。如果未配置日志传输,请按照以下步骤查询复制器的状态主题以获取错误消息。

如果ReplicatorFailure指标报告的值为非零值,请按照以下步骤进行故障排除:

  1. 配置一个可以连接到目标 MSK 集群并设置 Apache Kafka CLI 工具的客户端。请参阅连接到预置 Amazon MSK 集群。

  2. 在?处打开亚马逊 MSK 控制台 https://console.aws.amazon.com/msk/homeregion=us -east-1#/home/。

    获取 MSK 复制器和目标 MSK 集群的 ARN,并获取目标 MSK 集群的代理终端节点。

  3. 导出 MSK Replicator ARN 和代理端点:

    export TARGET_CLUSTER_SERVER_STRING=<BootstrapServerString> export REPLICATOR_ARN=<ReplicatorARN> export CONSUMER_CONFIG_FILE=<ConsumerConfigFile>
  4. 在您的<path-to-your-kafka-installation>/bin目录中,将以下脚本另存为query-replicator-failure-message.sh。

    #!/bin/bash # Script: Query MSK Replicator Failure Message # Description: This script queries exceptions from AWS MSK Replicator status topics # It takes a replicator ARN and bootstrap server as input and searches for replicator exceptions # in the replicator's status topic, formatting and displaying them in a readable manner # # Required Arguments: # --replicator-arn: The ARN of the AWS MSK Replicator # --bootstrap-server: The Kafka bootstrap server to connect to # --consumer.config: Consumer config properties file # Usage Example: # ./query-replicator-failure-message.sh --replicator-arn <replicator-arn> --bootstrap-server <bootstrap-server> --consumer.config <consumer.config> print_usage() { echo "USAGE: $0 ./query-replicator-failure-message.sh --replicator-arn <replicator-arn> --bootstrap-server <bootstrap-server> --consumer.config <consumer.config>" echo "--replicator-arn <String: MSK Replicator ARN> REQUIRED: The ARN of AWS MSK Replicator." echo "--bootstrap-server <String: server to connect to> REQUIRED: The Kafka server to connect to." echo "--consumer.config <String: config file> REQUIRED: Consumer config properties file." exit 1 } # Initialize variables replicator_arn="" bootstrap_server="" consumer_config="" # Parse arguments while [[ $# -gt 0 ]]; do case "$1" in --replicator-arn) if [ -z "$2" ]; then echo "Error: --replicator-arn requires an argument." print_usage fi replicator_arn="$2"; shift 2 ;; --bootstrap-server) if [ -z "$2" ]; then echo "Error: --bootstrap-server requires an argument." print_usage fi bootstrap_server="$2"; shift 2 ;; --consumer.config) if [ -z "$2" ]; then echo "Error: --consumer.config requires an argument." print_usage fi consumer_config="$2"; shift 2 ;; *) echo "Unknown option: $1"; print_usage ;; esac done # Check for required arguments if [ -z "$replicator_arn" ] || [ -z "$bootstrap_server" ] || [ -z "$consumer_config" ]; then echo "Error: --replicator-arn, --bootstrap-server, and --consumer.config are required." print_usage fi # Extract replicator name and suffix from ARN replicator_arn_suffix=$(echo "$replicator_arn" | awk -F'/' '{print $NF}') replicator_name=$(echo "$replicator_arn" | awk -F'/' '{print $(NF-1)}') echo "Replicator name: $replicator_name" # List topics and find the status topic topics=$(./kafka-topics.sh --command-config client.properties --list --bootstrap-server "$bootstrap_server") status_topic_name="__amazon_msk_replicator_status_${replicator_name}_${replicator_arn_suffix}" # Check if the status topic exists if echo "$topics" | grep -Fq "$status_topic_name"; then echo "Found replicator status topic: '$status_topic_name'" ./kafka-console-consumer.sh --bootstrap-server "$bootstrap_server" --consumer.config "$consumer_config" --topic "$status_topic_name" --from-beginning | stdbuf -oL grep "Exception" | stdbuf -oL sed -n 's/.*Exception:\(.*\) Topic: \([^,]*\), Partition: \([^\]*\).*/ReplicatorException:\1 Topic: \2, Partition: \3/p' else echo "No topic matching the pattern '$status_topic_name' found." fi

    运行此脚本来查询 MSK 复制器故障消息:

    <path-to-your-kafka-installation>/bin/query-replicator-failure-message.sh --replicator-arn $REPLICATOR_ARN --bootstrap-server $TARGET_CLUSTER_SERVER_STRING --consumer.config $CONSUMER_CONFIG_FILE

    此脚本输出了所有错误及其异常消息,以及受影响的主题分区。由于该主题包含了所有历史失败消息,因此应使用最后一条消息开始调查。以下是失败消息的示例:

    ReplicatorException: The request included a message larger than the max message size the server will accept. Topic: test, Partition: 1
常见故障和解决方案

以下内容介绍常见的 MSK 复制器故障以及如何缓解这些故障。

消息大小大于 max.request.size

原因:单个邮件大小超过 10 MB(默认最大值)。

下面是此失败消息类型的示例。

ReplicatorException: The message is 20635370 bytes when serialized which is larger than 10485760, which is the value of the max.request.size configuration. Topic: test, Partition: 1

解决方案:缩小主题中的单个消息大小。如果不能,请按照要求提高限额的说明进行操作。

消息大小大于服务器能接受的最大消息大小

原因:消息大小超过目标集群的最大消息大小。

下面是此失败消息类型的示例。

ReplicatorException: The request included a message larger than the max message size the server will accept. Topic: test, Partition: 1

解决方案:增加目标集群或主题的max.message.bytes配置。参见 max.message.bytes。

时间戳超出范围

原因:消息时间戳超出了目标集群的允许范围。

下面是此失败消息类型的示例。

ReplicatorException: Timestamp 1730137653724 of message with offset 0 is out of range. The timestamp should be within [1730137892239, 1731347492239] Topic: test, Partition: 1

解决方案:更新目标集群的message.timestamp.before.max.ms配置。参见 messag e.timestamp.before.max.ms。

记录批处理过大

原因:记录批量大小超过了目标集群上为主题设置的分段大小。MSK 复制器的最大批处理大小为 1MB。

下面是此失败消息类型的示例。

ReplicatorException: The request included message batch larger than the configured segment size on the server. Topic: test, Partition: 1

解决方案:将目标集群更新segment.bytes为至少 1048576 (1 MB)。参见 segment.by tes。

注意

如果在应用这些解决方案后ReplicatorFailure指标继续发出非零值,请重复故障排除过程,直到该指标发出零值。

排除来自自管理的 Kafka 集群的复制问题

MSK Replicator 无法连接到自建的 Kafka 集群

如果 MSK Replicator 无法连接到您的自建的 Kafka 集群,请执行以下检查:

  1. 验证您的 VPN 或直接连接是否处于活动状态以及路由表是否正确。

  2. 验证安全组是否允许来自 MSK Replicator 子网的入站流量通过 SASL_SSL 端口(通常为 9096)。

  3. 验证从 VPC 到自管理集群代理主机名的 DNS 解析。

  4. 检查亚马逊中的KafkaClusterPingSuccessCount指标 CloudWatch ——值为 0 表示连接故障。

SASL/SCRAM 或 mTLS 身份验证失败

如果AuthError指标为非零或复制器日志显示 SASL/SCRAM 或 mTLS 错误:

  1. 验证存储在 S AWS ecrets Manager 中的凭据是否与自管集群上的 SCRAM 用户凭证相匹配(用于 SASL/SCRAM),或者客户端证书和私钥是否正确(对于 mTL)。

  2. 验证 SCRAM 用户或证书主体是否具有所需的 ACL 权限(读取、描述主题;读取、描述消费群体;在集群上描述)。

  3. 检查AuthError指标以确认身份验证错误,并使用该ClusterAlias维度确定源集群或目标集群是否受到影响。

SASL/OAUTHBEARER (OAuth) 身份验证失败

如果AuthError指标为非零或复制器日志显示访问令牌检索或 SASL/OAUTHBEARER 错误:

  1. 验证是否tokenEndpointUrl正确,使用 HTTPS 架构,并且可以从您为复制器提供的 VPC 子网进行访问。A KafkaClusterPingSuccessCount 为 0 加上令牌错误通常表示无法访问令牌端点。

  2. 对于客户端凭证机制,请验证 S AWS ecrets Manager client_secret 中的client_id和是否正确,以及令牌获取机制是否符合您的 IDP 期望。

  3. 验证对您的代币获取机制tokenEndpointAuthenticationMethod是否有效。客户端凭证机制需要POST或BASIC,而客户端凭证断言机制需要NONE。

  4. 对于 IAM JWT 持有者和客户端凭证断言机制,请验证服务执行角色是否具有sts:GetWebIdentityToken权限,audience且signingAlgorithm与您的 IDP 验证令牌时的期望相匹配。

  5. 如果您的 IDP 使用私有 CA,请验证所引用的 CA 证书tokenEndpointTlsCertificateArn是否完整且有效。

  6. 验证您的 IDP 映射访问令牌的 Kafka 主体是否具有 MSK Replicator 在源集群上要求的 ACL 权限。

  7. 检查复制器日志中是否有令牌端点返回的 HTTP 状态和 OAuth error 代码。MSK Replicator 记录这些字段,但从不记录令牌端点响应正文。

SSL 证书问题

如果复制器无法与自管集群建立安全连接:

  1. 验证 AWS Secrets Manager 中的certificate值是否包含 PEM 格式的完整 CA 证书链。

  2. 验证是否在所有自托管集群代理上配置了 SSL 侦听器。

  3. 确认证书未过期且由可信 CA 颁发。

自建集群的消费组抵消同步失败

如果消费组偏移量未正确同步:

  1. 验证ConsumerGroupOffsetSyncFailure指标 — 它应该是 0。

  2. 确认使用者组正在源集群上积极消费(非活动使用者组可能无法同步)。

  3. 对于双向复制,请验证两个复制器true上均设置synchroniseConsumerGroupOffsets为。