【问题标题】:Cuba Platform - Spring Kafka Integration古巴平台 - Spring Kafka 集成
【发布时间】:2020-07-29 10:10:19
【问题描述】:

我必须将 Kafka 集成到 Cuba,我认为这就像添加 spring kafka 依赖项并创建一个 Configuration 注释类来初始化 Kafka Consumer 一样简单,因为 Cuba 是基于 Spring 的。

当我添加一个配置时,我发现启动古巴时没有扫描它。当我切换到 CUBA 视图时,我注意到只有那些注释为 Service 或 Component 的类会被读取。但是,即使我添加了一个Component 类,它仍然没有被正确扫描(我添加了一个用@Value 注释的字段,它查找一个不存在的属性,但古巴在我启动它时没有抛出任何错误)

【问题讨论】:

    标签: apache-kafka spring-kafka cuba-platform


    【解决方案1】:

    有一个关于CUBA+Kafka集成的简单例子,你可以在这里找到:https://github.com/cuba-labs/kafka-sample

    配置过程取自official Spring documentation。

    1. 关键配置类是com.company.kafkasample.config.KafkaConfig。它包含许多可帮助您配置 Kafka 设施的 bean。在此特定示例中,同时配置了生产者和消费者。请注意,配置参数是硬编码的,但这只是一个示例。
        @Bean
        public ConsumerFactory<Integer, String> consumerFactory() {
            return new DefaultKafkaConsumerFactory<>(consumerConfigs());
        }
    
        @Bean
        public Map<String, Object> consumerConfigs() {
            Map<String, Object> props = new HashMap<>();
            props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
            props.put(ConsumerConfig.GROUP_ID_CONFIG, "sample-kafka");
            props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
            props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
            props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
            props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
            props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            return props;
        }
    
    1. 之后,您应该能够通过注入 KafkaTemplate bean 创建一个可以向 Kafka 队列发送消息的服务。
        @Inject
        private KafkaTemplate<Integer, String> template;
    
    
        @Override
        public void sendMessage(String message) {
            log.info("Sending {} using Kafka", message);
            long id = uniqueNumbersService.getNextNumber("users");
            ListenableFuture<SendResult<Integer, String>> send = template.send("users", (int) id, message);
            send.addCallback(new ListenableFutureCallback<SendResult<Integer, String>>() {
                @Override
                public void onFailure(Throwable ex) {
                   log.info("Failed to send message {}, error {}", message, ex.getMessage());
                }
    
                @Override
                public void onSuccess(SendResult<Integer, String> result) {
                    log.info("Message {} sent", message);
                }
            });
        }
    
    1. 然后您可以将此服务注入屏幕并在那里使用它。
    2. 对于receiver,你可以在CUBA组件中使用@KafkaListener注解来注解它的方法。例如,下面的示例将 kafka 消息保存到数据库中。
    @Component
    @DependsOn("consumerFactory")
    public class MessageListener {
    
        @Inject
        private DataManager dataManager;
    
        @KafkaListener(id = "sample-kafka", topics = "users")
        public void listen1(String foo, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) int id) {
            KafkaMessage kafkaMessage = dataManager.create(KafkaMessage.class);
            kafkaMessage.setKafkaID(id);
            kafkaMessage.setContent(foo);
            dataManager.commit(kafkaMessage);
        }
    

    【讨论】:

    • 虽然这在理论上可以回答问题,it would be preferable 在这里包含答案的基本部分,并提供链接以供参考。
    猜你喜欢
    • 1970-01-01
    • 2023-03-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-31
    • 1970-01-01
    相关资源
    最近更新 更多