You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ECS Fargate中Go应用无法连接MSK Serverless集群问题排查

问题:ECS Fargate上的Go Sarama消费者无法解析MSK Serverless引导代理域名

我有一个用Go编写的消费者应用,部署在ECS Fargate上,尝试从MSK Serverless集群消费消息。ECS服务与MSK集群处于同一VPC,MSK关联的安全组已更新为允许来自ECS服务关联安全组的入站连接。我在EC2上还有一个Node.js应用,能够成功向该集群生产消息。然而,ECS应用中出现了来自Sarama的如下日志:

Failed to connect to broker 'boot-..kafka-serverless.us-east-2.amazonaws.com:9098: dial tcp: lookup 'boot-..kafka-serverless.us-east-2.amazonaws.com: no such host

集群的引导代理地址是正确的,查看Sarama源码后发现,该错误发生在Bearer认证逻辑之前。我也已确认ECS服务所关联的角色拥有连接Kafka并消费消息的正确权限。

以下是连接集群的相关源码:

import (
    "context"
    "os"

    "github.com/IBM/sarama"
    "github.com/aws/aws-msk-iam-sasl-signer-go/signer"
    "github.com/<redacted>/log"
)

type MSKConf struct {
    Region string `json:"region"`
}

func (m *MSKConf) Token() (*sarama.AccessToken, error) {
    signer.AwsDebugCreds = true
    token, _, err := signer.GenerateAuthToken(context.Background(), m.Region)
    return &sarama.AccessToken{Token: token}, err
}

type ConsumerConfig struct {
    Brokers  []string `json:"brokers"`
    Topic    string   `json:"topic"`
    ClientID string   `json:"clientID"`
    MSK      *MSKConf `json:"msk"`
}

func (cc *ConsumerConfig) toSaramaConf(ctx context.Context) *sarama.Config {
    config := sarama.NewConfig()
    config.ClientID = cc.ClientID
    if cc.MSK != nil {
        log.Info("Connecting to MSK...")
        config.Net.SASL.Enable = true
        config.Net.SASL.Mechanism = sarama.SASLTypeOAuth
        config.Net.SASL.TokenProvider = cc.MSK
    }
    return config
}

func InitConsumer(ctx context.Context, conf *ConsumerConfig, handler IHandler) error {
    sarama.Logger = log.New(os.Stdout, "[sarama] ", log.LstdFlags)
    consumer, err := setUpConsumer(ctx, conf)
    if err != nil {
        return err
    } else {
        log.Infof("Kafka Consumer is up and running!")
    }
    err = consumeMessages(ctx, handler, consumer, conf)
    if err != nil {
        return err
    }
    return nil
}

func setUpConsumer(ctx context.Context, conf *ConsumerConfig) (sarama.Consumer, error) {
    config := conf.toSaramaConf(ctx)
    return sarama.NewConsumer(conf.Brokers, config)
}

解决方案

核心问题是ECS Fargate任务无法解析MSK引导代理域名,按以下步骤排查解决:

  • 修正引导代理地址格式错误:日志中显示的域名boot-..kafka-serverless.us-east-2.amazonaws.com存在连续的两个点,明显是配置错误。重新从MSK控制台复制正确的引导代理地址,检查ConsumerConfig的Brokers字段是否存在拼写或格式问题。
  • 检查VPC子网DNS配置:确保Fargate任务所在的VPC子网开启了EnableDnsHostnames和EnableDnsSupport属性。MSK Serverless的域名依赖VPC内部DNS解析,这两个属性未开启会导致解析失败。同时确认子网路由表包含指向VPC DNS服务器(VPC CIDR段+2)的路由。
  • 验证网络ACL和出站权限:检查子网的网络ACL是否允许UDP 53端口的出站流量(DNS解析需要),同时确保Fargate任务的安全组允许出站访问VPC DNS服务器。
  • 补充Sarama TLS配置:MSK Serverless强制要求TLS连接,在toSaramaConf方法中添加TLS开启配置:
    config.Net.TLS.Enable = true
    
  • 测试域名解析:在Fargate任务中添加临时命令nslookup <正确的引导代理域名>,手动验证域名是否能正常解析,快速定位是配置问题还是网络问题。

内容的提问来源于stack exchange,提问作者user2704305

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 05:25:16