【问题标题】:how to consume events from kafka by a Spring Rest endpoint如何通过 Spring Rest 端点使用来自 kafka 的事件
【发布时间】:2021-05-06 20:02:44
【问题描述】:

我是卡夫卡的新手。我已经看到消费者“一直在运行”,并在发布主题后立即从主题中检索消息。

在典型的数据库 Web 应用程序中,您有一个连接到 DB 并返回一些响应的 rest API。

据我所见,消费者保持活跃且永不关闭。 所以我不知道如何根据客户端请求从主题返回消息子集。

我认为该服务会创建一个消费者来获得我需要的东西,但就消费者永远不会关闭而言,我想我的意见是不正确的。 我该怎么办?

【问题讨论】:

  • 你的意思是org.springframework.kafka.annotation.KafkaListener 不适合你吗?
  • @k-wasilewski 我猜 KafkaListener 仍在等待接收来自某个主题的新消息。我希望我的客户调用具有开始时间和结束时间间隔的 API。该服务应接收该时间间隔内的所有消息。

标签: spring apache-kafka kafka-consumer-api


【解决方案1】:

然后这是一个简单的问题,即持久化通过 KafkaListener 接收到的消息,假设将它们中的每一个添加到一个简单的集合(连同它的时间戳)并实现一个端点以相应地过滤消息并返回其中的一些。

@Controller
public class KafkaController {
    @Autowired
    private KafkaProducerConfig kafkaProducerConfig;
    private Map<Date, String> msgMap = new HashMap();

    @KafkaListener(topics = "myTopic", groupId = "myGroup")
    public void listenAndAddMsg(String message) {
        msgMap.put(new Date(), message);
    }

    @PostMapping("messages")
    @ResponseBody
    public String filterMessages(@RequestBody Interval interval) {
        return msgMap.entrySet() 
          .stream() 
          .filter(map -> map.getKey().after(interval.getStartDate()) && map.getKey().before(interval.getEndDate()))
          .collect(Collectors.toMap(map -> map.getKey(), map -> map.getValue()));  
    }
}

public class Interval {
    private Date startDate;
    private Date endDate;

    // setters and getters
}

【讨论】:

  • 在您的示例中,我猜您正在将每条消息添加到地图中,并且您只需按该地图进行过滤。这意味着将所有内容都存储在内存中。我可以将其直接保存到内存中,而不是使用 Kafka,但我每天应该有 1.500.000 个对象并最多存储 3 个月......我想告诉消费者:给我从这个日期到这个日期的所有消息(通常不超过 3 天)。有可能吗?
  • 不,这是您作为后端开发人员关心的问题。因此,与其将消息存储在地图中,不如将它们存储在数据库中......
猜你喜欢
  • 1970-01-01
  • 2023-03-18
  • 1970-01-01
  • 2018-03-17
  • 2017-06-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多