Spring Integration Java DSL集成Kafka:数据库记录无法发送至Topic求助
问题分析与解决方案
核心问题
你的代码中Kafka出站适配器的配置方式错误,导致消息没有被实际发送到Kafka。在outboundChannelAdapterFlow的handle方法中,你通过lambda表达式创建了Kafka.outboundChannelAdapter实例,但并没有调用它处理消息,只是返回了适配器对象,因此消息无法投递到Kafka集群。
修正步骤
1. 修复Kafka出站适配器的配置逻辑
将handle方法中的lambda替换为直接配置Kafka.outboundChannelAdapter作为消息处理器,让Spring Integration自动调用适配器完成消息发送。
2. 修正ProducerFactory泛型类型
你的ProducerFactory泛型定义为<Integer, String>,但消息key使用的是String类型,会导致类型不匹配,需调整为<String, String>。
3. 增加Header空值保护
避免topic、key等Header为空时触发空指针异常,建议通过SpEL表达式添加默认值。
修正后的完整代码
@Configuration public class KafkaProduceConfig { @Bean public IntegrationFlow pollingAdapterFlow(EntityManagerFactory entityManagerFactory, MyTransformer transformer) { return IntegrationFlow .from(Jpa.inboundAdapter(entityManagerFactory).entityClass(MyRecord.class), e -> e.poller(p -> p.cron("*/1 * * * * *").maxMessagesPerPoll(1).transactional()) .autoStartup(true)) .log(message -> "Polled DB Records from KafkaProduceConfig : " + message.getPayload()) .split() .log(message -> "Record after split : " + message.getPayload()) .enrichHeaders(hrdSpec ->hrdSpec.headerExpression("myRecord", "payload",true)) .transform(transformer,"getCustomeRecord") .enrichHeaders(hrdSpec ->hrdSpec.headerExpression("customeRecord","payload",true)) .log(message -> "Transformed Record : " + message.getPayload() +",topic :"+message.getHeaders().get("topic")) .channel("sendToKafka") .get(); } @Bean public IntegrationFlow outboundChannelAdapterFlow() { return IntegrationFlow.from("sendToKafka") .log(message -> "outboundChannelAdapterFlow received payload : " + message.getPayload() +",topic :" +message.getHeaders().get("topic")+"key :"+message.getHeaders().get("key")) // 修正:直接将Kafka出站适配器作为handle的参数 .handle(Kafka.outboundChannelAdapter(producerFactory()) .topicExpression("headers['topic'] ?: 'default-topic'") .messageKeyExpression("headers['key'] ?: 'default-key'") .partitionIdExpression("headers['partitionId'] ?: 0")) .get(); } // 修正泛型,匹配key的String类型 public ProducerFactory<String, String> producerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 可选:添加acks配置确保消息被Kafka确认 props.put(ProducerConfig.ACKS_CONFIG, "all"); return new DefaultKafkaProducerFactory<>(props); } }
额外排查建议
- 开启Kafka生产者调试日志,查看发送细节:在
application.properties中添加logging.level.org.springframework.kafka=DEBUG - 确认Kafka集群正常运行,
localhost:9092可访问,目标Topic已存在(或开启自动创建:props.put(AdminClientConfig.AUTO_CREATE_TOPICS_CONFIG, true)) - 检查
MyTransformer.getCustomeRecord方法返回的payload是否为String类型,匹配配置的StringSerializer
内容的提问来源于stack exchange,提问作者user23543923
相关产品推荐
相关产品推荐

