【问题标题】:Spring-Integration XML to JavaSpring 集成 XML 到 Java
【发布时间】:2015-07-22 16:31:39
【问题描述】:

如何将此代码转换为 Java 配置?

<int-kafka:outbound-channel-adapter
        id="mainOutboundChannelAdapter"
        kafka-producer-context-ref="kafkaProducerContext"
        channel="mainOutboundTopicChanel">
</int-kafka:outbound-channel-adapter>

【问题讨论】:

    标签: java spring spring-integration


    【解决方案1】:

    是的,你可以。请找到最新的Spring Integration Java DSL

    您的情况可能如下所示:

    @Bean
    public IntegrationFlow sendToKafkaFlow(String serverAddress) {
        return f -> f.<String>split(p -> FastList.newWithNValues(100, () -> p), null)
                .handle(kafkaMessageHandler(serverAddress));
    }
    
    private KafkaProducerMessageHandlerSpec kafkaMessageHandler(String serverAddress) {
        return Kafka.outboundChannelAdapter(props -> props.put("queue.buffering.max.ms", "15000"))
                .messageKey(m -> m.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER))
                .addProducer(TEST_TOPIC, serverAddress, this::producer);
    }
    
    private void producer(KafkaProducerMessageHandlerSpec.ProducerMetadataSpec metadata) {
        metadata.async(true)
                .batchNumMessages(10)
                .valueClassType(String.class)
                .<String>valueEncoder(String::getBytes)
                .keyEncoder(new IntEncoder(null));
    }
    

    更新 没有 Lambda,但仍然是 Spring 集成:

    @Bean
    @ServiceActivator(inputChannel = "mainOutboundTopicChanel")
    public MessageHandler kafkaProducer() {
        return new KafkaProducerMessageHandler<String, String>(kafkaProducerContext());
    }
    
    @Bean
    public KafkaProducerContext<String, String> kafkaProducerContext() {
        KafkaProducerContext<String, String> kafkaProducerContext = new KafkaProducerContext<String, String>();
        ProducerMetadata<String, String> producerMetadata = new ProducerMetadata<String, String>(TOPIC);
        producerMetadata.setValueClassType(String.class);
        producerMetadata.setKeyClassType(String.class);
        Encoder<String> encoder = new StringEncoder<String>();
        producerMetadata.setValueEncoder(encoder);
        producerMetadata.setKeyEncoder(encoder);
        producerMetadata.setAsync(true);
        Properties props = new Properties();
        props.put("queue.buffering.max.ms", "15000");
        ProducerFactoryBean<String, String> producer =
                new ProducerFactoryBean<String, String>(producerMetadata, kafkaRule.getBrokersAsString(), props);
        ProducerConfiguration<String, String> config =
                new ProducerConfiguration<String, String>(producerMetadata, producer.getObject());
            kafkaProducerContext.setProducerConfigurations(Collections.singletonMap(TOPIC, config));
        return kafkaProducerContext;
    }
    

    别忘了在@Configuration 旁边添加@EnableIntegration

    未来:Spring 中的任何 XML 标记都会被一些 NamespaceHandler 解析,例如在这种情况下,它是KafkaNamespaceHandler。阅读其源代码,我们可以找到以下几行:

    registerBeanDefinitionParser("outbound-channel-adapter", new KafkaOutboundChannelAdapterParser());
            registerBeanDefinitionParser("producer-context", new KafkaProducerContextParser());
    

    当我们转到KafkaOutboundChannelAdapterParser 并看到它填充了BeanDefinition

    final BeanDefinitionBuilder kafkaProducerMessageHandlerBuilder =
                                    BeanDefinitionBuilder.genericBeanDefinition(KafkaProducerMessageHandler.class);
    

    源代码等等。

    更新 2

    Consumer 部分:

    @Bean
    @InboundChannelAdapter(value = "fromKafkaChannel",
        poller = @Poller(fixedRate = "10", maxMessagesPerPoll = "1"))
    public MessageSource<Map<String, Map<Integer, List<Object>>>> kafkaMessageSource() {
        return new KafkaHighLevelConsumerMessageSource<String, String>();
    }
    
    @Bean
    public KafkaConsumerContext<String, String> kafkaConsumerContext() {
        KafkaConsumerContext<String, String> kafkaConsumerContext = new KafkaConsumerContext<String, String>();
        .....
        kafkaConsumerContext.setConsumerConfigurations(map);
        return kafkaConsumerContext;
    }
    

    【讨论】:

    • 嗨,Artem,感谢您的回复,我以前见过这类答案。我觉得人们因为 java 8 lambda 特性而远离这些。常规 java 可能对许多还不了解 java 8 的人更有帮助。
    • 是的...我理解您(和其他人)的担忧...请在我的回答中找到更新。
    • 感谢您,作为一名初级开发人员,我仍在学习如何深入研究不同的场景等等。您的代码非常有用。如果问得不算多,您是否会碰巧将相应的 ConsumerContext 与 java config 一起发布?我似乎在谷歌上找不到稳定的教程。
    • 请在我的回答中找到@InboundChannelAdapter 的更新。并随时为 Java Config 文档提出 GitHub 问题 (github.com/spring-projects/spring-integration-kafka/issues)。
    • 伙计,非常感谢。你不知道我被困在这上面多久了。谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-08-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-16
    相关资源
    最近更新 更多