Java创建Kafka Topic报错:scala.Product.$init$方法不存在
我最近基于Kafka 2.12和Java 1.8开发动态创建Topic的功能,写了一个接收Topic名称作为参数的创建方法,但运行时遇到了NoSuchMethodError,折腾半天终于搞明白问题所在,分享给大家:
我的实现代码
private static void CreateKafkaTopic(String topicName) { ZkClient zkClient = null; ZkUtils zkUtils = null; try { String zookeeperConnect = "localhost:2181"; int sessionTimeOutInMs = 15 * 1000; // 15 secs int connectionTimeOutInMs = 10 * 1000; // 10 secs zkClient = new ZkClient(zookeeperConnect, sessionTimeOutInMs, connectionTimeOutInMs, ZKStringSerializer$.MODULE$); boolean isSecureKafkaCluster = false; zkUtils = new ZkUtils(zkClient, new ZkConnection(zookeeperConnect), isSecureKafkaCluster); Properties topicConfig = new Properties(); AdminUtils.createTopic(zkUtils, topicName, 1, 1, topicConfig,RackAwareMode.Disabled$.MODULE$); } catch (Exception ex) { ex.printStackTrace(); } finally { if (zkClient != null) { zkClient.close(); } } }
遇到的错误
Exception in thread "main" java.lang.NoSuchMethodError: scala.Product.$init$(Lscala/Product;)V at kafka.admin.RackAwareMode$Disabled$.(RackAwareMode.scala:27) at kafka.admin.RackAwareMode$Disabled$.(RackAwareMode.scala) at com.OTMProducer.CreateKafkaTopic(OTMProducer.java:243)
问题原因分析
这个错误本质是Scala版本不兼容导致的。你用的Kafka 2.12对应Scala 2.12版本,但项目里可能引入了其他版本的Scala依赖(比如Scala 2.11),或者Kafka客户端依赖的Scala版本和项目里的不一致。scala.Product.$init$这个方法在不同Scala版本里的签名有变化,当JVM加载到不匹配的类时,就会抛出这个找不到方法的错误。
另外要注意,从Kafka 2.0开始,官方已经不推荐用AdminUtils这种依赖ZooKeeper的方式创建Topic了,更建议使用Kafka AdminClient API——这种方式不需要直接操作ZK,更符合Kafka的架构设计,也能避免很多依赖兼容问题。
解决方案
方案1:修复Scala版本依赖
确保项目中所有Scala相关的依赖(包括Kafka客户端依赖)都统一使用和Kafka版本匹配的Scala版本。Kafka 2.12对应的Scala版本是2.12.x,在Maven或Gradle中明确指定版本,排除掉冲突的依赖。
比如Maven中可以这样配置:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka_2.12</artifactId> <version>你的Kafka版本号</version> </dependency>
注意artifactId里的_2.12就是指定Scala版本为2.12,一定要和你的Kafka版本对应。
方案2:改用AdminClient API(推荐)
直接用官方推荐的AdminClient来创建Topic,代码更简洁,也不会有ZK依赖的兼容问题。示例代码如下:
private static void createKafkaTopic(String topicName) { Properties adminProps = new Properties(); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); // Kafka broker地址 adminProps.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, "10000"); try (AdminClient adminClient = AdminClient.create(adminProps)) { // 检查Topic是否已存在 ListTopicsResult topicsResult = adminClient.listTopics(); Set<String> existingTopics = topicsResult.names().get(); if (existingTopics.contains(topicName)) { System.out.println("Topic " + topicName + " already exists"); return; } // 创建Topic的配置 NewTopic newTopic = new NewTopic(topicName, 1, (short) 1); // 分区数1,副本数1 // 如果需要额外配置,可以添加 // Map<String, String> configs = new HashMap<>(); // configs.put("cleanup.policy", "delete"); // newTopic.configs(configs); CreateTopicsResult result = adminClient.createTopics(Collections.singleton(newTopic)); result.all().get(); // 等待创建完成 System.out.println("Topic " + topicName + " created successfully"); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); } }
这个方法只需要依赖Kafka的客户端包,不需要引入ZooKeeper相关的依赖,也避开了Scala版本兼容的坑,更适合新版本的Kafka。
内容的提问来源于stack exchange,提问作者user1708054

