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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:25:28