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

Java Lambda同VPC连接MSK创建Topic抛出TimeoutException

问题现象

在与MSK集群同VPC、同安全组下创建的Java Lambda函数执行时,CloudWatch日志抛出如下异常:

org.apache.kafka.common.errors.TimeoutException

用于创建Kafka Topic的Java代码如下:

public String handleRequest(SQSEvent input, Context context) {

        LambdaLogger logger = context.getLogger();
        if(bootStrapServer == null) {
            System.out.println("missing boot strap server env var");
            return "Error, bootStrapServer env var missing";
        }

        Properties props = new Properties();
        props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServer);
        props.put(AdminClientConfig.CLIENT_ID_CONFIG, "java-data-screaming-demo-lambda");
        props.put(AdminClientConfig.SECURITY_PROTOCOL_CONFIG, "PLAINTEXT");

        try {
            this.createTopic("TestLambdaTopic", props, logger);
        } catch (Exception e) {
            logger.log("err in creating topic: " + gson.toJson(e));
        }


        return "Ok";
    }

    public void createTopic(String topicName, Properties properties, LambdaLogger logger ) throws Exception {
        try (Admin admin = Admin.create(properties)) {
            int partitions = 1;
            short replicationFactor = 2;
            NewTopic newTopic = new NewTopic(topicName, partitions, replicationFactor);

            List<NewTopic> topics = new ArrayList<NewTopic>();
            topics.add(newTopic);

            CreateTopicsResult result = admin.createTopics(topics);
            // get the async result for the new topic creation
            KafkaFuture<Void> future = result.values().get(topicName);

            // call get() to block until topic creation has completed or failed
            future.get();
            if (future.isDone()) {
                logger.log("future is done");
            }

            logger.log("what is result from create topics: " + gson.toJson(result));
        }
    }
排查方案与修复方法

按优先级从高到低排查以下问题:

  • Bootstrap Server地址配置错误
    不要使用ZooKeeper连接串作为Kafka bootstrap地址,必须填写MSK集群明文监听器(对应9092端口)的broker端点列表,多个地址用英文逗号分隔。如果你的MSK集群开启了IAM/SASL认证,PLAINTEXT安全协议无法连通,需要替换为对应认证协议配置。
  • 安全组规则未放通访问
    即使Lambda和MSK使用同一个安全组,也需要在该安全组的入站规则中,添加来源为安全组自身、端口为9092的TCP允许规则,否则同安全组内的资源默认无法互通。同时检查安全组出站规则,确认没有拦截Lambda到9092端口的出站请求。
  • 副本因子配置与集群规模不匹配
    代码中写死的副本因子为2,如果你的MSK是单节点开发集群(broker数量为1),创建Topic的请求会因为副本数无法满足一直阻塞,最终触发超时。单节点集群需要把replicationFactor参数改为1。
  • 网络路由/NACL拦截
    检查Lambda关联子网的路由表,确认到MSK broker所在网段的路由没有被黑洞拦截;同时检查子网关联的网络ACL,确认没有拒绝9092端口的入站、出站流量。如果Lambda部署在公有子网,需要给Lambda的弹性网卡分配公网IP,否则会出现内网访问异常。
  • 客户端超时时间过短
    Lambda冷启动阶段网络初始化耗时较长,Kafka AdminClient默认的请求超时时间可能不足,可以在配置中增加超时参数:
// 请求超时调整为60秒
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 60000);
// API调用总超时调整为60秒
props.put(AdminClientConfig.DEFAULT_API_TIMEOUT_MS_CONFIG, 60000);
// 增加重试次数
props.put(AdminClientConfig.RETRIES_CONFIG, 3);

排查时可以先在Lambda代码中增加Socket连通性测试,启动时先尝试连接bootstrap地址的9092端口,确认网络层连通后再排查Kafka侧配置问题。

内容的提问来源于stack exchange,提问作者Will Zefeng Qiu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 22:03:29