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
相关产品推荐
相关产品推荐

