Kafka Producer传入字符串Key发送 消费读取为null问题排查
问题:Kafka Producer显式传入非空Key,消费端读取到Key为null
我在使用Producer向Kafka主题发送消息,开展JUnit测试时发现:应用代码中定义的Producer发送的消息Key被识别为null,但JUnit测试类内部定义的Producer发送同逻辑消息时Key可以被正常识别,我已经明确为消息传入了字符串类型的Key。
相关代码
主应用类代码
final Producer<String, HashSet<String>> actualApplicationProducer; ApplicationInstance(String bootstrapServers) // 构造方法 { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.CLIENT_ID_CONFIG, "ActualClient"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, CustomSerializer.class.getName()); props.put(ProducerConfig.LINGER_MS_CONFIG, lingerBatchMS); props.put(ProducerConfig.BATCH_SIZE_CONFIG, Math.min(maxBatchSizeBytes,1000000)); actualApplicationProducer = new KafkaProducer<>(props); } public void doStuff() { HashSet<String> values = new HashSet<String>(); String key = "applicationKey"; // 该行代码发送的消息Key被识别为null actualApplicationProducer.send(new ProducerRecord<>(topicName, key, values)); }
JUnit测试类代码
@EmbeddedKafka @ExtendWith(SpringExtension.class) @SuppressWarnings("static-method") @TestInstance(TestInstance.Lifecycle.PER_CLASS) public class CIFFileProcessorTests { /** 单元测试用嵌入式Kafka Broker */ @Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; @BeforeAll public void setUpBeforeClass(@TempDir File globalTablesDir, @TempDir File rootDir) throws Exception { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.CLIENT_ID_CONFIG, "JUnitClient"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, CustomSerializer.class.getName()); props.put(ProducerConfig.LINGER_MS_CONFIG, lingerBatchMS); props.put(ProducerConfig.BATCH_SIZE_CONFIG, Math.min(maxBatchSizeBytes,1000000)); try(Producer<String, HashSet<String>> junitProducer = new KafkaProducer<>(props)) { HashSet<String> values = new HashSet<>(); // 此处和主应用代码逻辑完全一致,发送的Key可被正常识别,不为null junitProducer.send(new ProducerRecord<>(topicName,"junitKey",values)); } } @Test public void test() { ApplicationInstance sut = new ApplicationInstance(embeddedKafkaBroker.getBrokersAsString()); sut.doStuff(); // records为LinkedBlockingQueue,由监听主题的KafkaMessageListenerContainer通过MessageListener填充 ConsumerRecord<String, HashSet<String>> record = records.poll(1,TimeUnit.SECONDS); assertEquals("junitKey", record.key()); // 测试失败:实际返回null } }
自定义序列化器代码
try (final ByteArrayOutputStream baos = new ByteArrayOutputStream(); ObjectOutputStream oos = new ObjectOutputStream(baos)) { oos.writeObject(object); return baos.toByteArray(); }
排查更新
我检查发送返回的元数据后确认,Producer实际确实发送了非null的Key:
RecordMetadata info = actualApplicationProducer.send(new ProducerRecord<>(topicName, key, values)).get(); System.out.println("INFO - partition: " + info.partition() + ", topic: " + info.topic() + ", offset: " + info.offset() + ", timestamp: "+ info.timestamp() + ", keysize: " + info.serializedKeySize() + ", valuesize: " + info.serializedValueSize());
输出结果:
INFO - partition: 0, topic: topicName, offset: 2, timestamp: 1656060840304, keysize: 14, valuesize: 6258
序列化后Key大小大于0,证明Producer端确实没有发送null Key,问题出在消费读取环节。
问题原因&解决方案
这问题根本不是Producer发了null Key,是测试环节几个常见坑叠在一起了:
- 消费端反序列化配置错配:你写的自定义序列化器是用
ObjectOutputStream写Java对象的,这个序列化器只配给了Value使用,Key是用StringSerializer序列化的纯字符串字节。你给消费端监听容器配反序列化器的时候,大概率把Key的反序列化器也配成了读Java对象的自定义反序列化器,用ObjectInputStream去读String序列化出来的字节,要么读失败被监听容器吞了异常返回null,要么直接读出来乱码转成null。 - 批量消息没刷盘就开始拉取:你给Producer配了
linger.ms和批量发送参数,调用send()之后消息会先存在本地内存缓存里凑批次,不会立刻发到Broker。你调用完doStuff()直接poll,根本没等消息发出去,自然拿不到目标消息。 - 拉取逻辑混乱:你在
@BeforeAll里还先发了一条Key为junitKey的消息,poll是按偏移量从早到晚拿的,就算拉到消息也是先拿到之前那条旧的,你还把断言预期值写成junitKey,逻辑本身就有问题。
修复步骤
- 先改消费端配置:Key的反序列化器明确指定
StringDeserializer,Value才用你自己写的对应HashSet的自定义反序列化器,两者不要搞混。 - 发完消息强制刷盘:在
sut.doStuff()后面加一行actualApplicationProducer.flush(),把内存里攒的批量消息立刻推到Broker,不要等linger时间到。 - 调整拉取断言逻辑:要么测试前把主题里的旧数据全清空,要么循环poll几次,直到拿到对应Key的消息再做断言,避免被之前测试残留的旧消息干扰。
内容的提问来源于stack exchange,提问作者simonalexander2005
相关产品推荐
相关产品推荐

