【问题标题】:publish multiple events shares some attributes in one kafka topic发布多个事件在一个kafka主题中共享一些属性
【发布时间】:2022-11-18 04:28:36
【问题描述】:

我需要从代表员工旅程事件的同一个项目发布多条消息,并且我只需要使用一个主题来发布这些消息,因为它们代表同一个项目,但在某些情况下,消息可能包含额外的字段,例如:

所有消息共享(id、名称、类型、日期) 有时某些事件可能会有更多字段,例如(课程 ID、课程名称)

所以我打算使用一个名为“Journey”的父对象,其中包含“Event”对象, 如果需要,我将创建多个子对象,如扩展此事件的 LMSEvent 等,并使用 Jackson + spring boot over rest API 来根据类型属性执行所需的转换,然后直接将此消息发布到 Kafka,因此,每个对象包含它自己的属性。

对于消费者,我会做一些策略模式,并在需要时为每种类型做所需的逻辑。

消息大小不会很大,我不希望每个事件有更多不同的属性。

我想知道这种方法是否好,如果不好,还有什么选择。

多谢

【问题讨论】:

    标签: spring-boot apache-kafka


    【解决方案1】:

    我认为总的来说这是个好方法。关于主题或多个模式的单一消息模式总是一个好问题,两者都有一些亮点和缺点,您可以在Martin Kleppmann article.中阅读更多相关信息

    当您决定在单个主题上有多个事件时,从 rest api 开始,然后是 Kafka 生产者和消费者,您可以使用相同的方法序列化和反序列化事件,@JsonTypeInfo@JsonSubTypes 可以完成这项工作:

    @JsonTypeInfo(
            use = JsonTypeInfo.Id.NAME,
            include = JsonTypeInfo.As.EXISTING_PROPERTY,
            property = "type")
    @JsonSubTypes({
            @JsonSubTypes.Type(value = LMSEvent.class, name = "LMSEvent"),
            @JsonSubTypes.Type(value = YetAnotherEvent.class, name = "YetAnotherEvent")
    })
    public interface Event {
        String getType();
        default boolean hasType(String type) {
            return getType().equalsIgnoreCase(type);
        }
        default <T> T getConcreteEvent(Class<T> clazz) {
            return clazz.cast(this);
        }
    }
    

    当您使用 spring-kafka 使用该类型的消息时,您可以定义一些非常简洁的代码,其中每个方法都使用具体的事件类型,因此您不需要自己编写一些肮脏的转换:

        @KafkaListener(topics = "someEvents", containerFactory = "myKafkaContainerFactory")
        public class MyKafkaHandler {
    
            @KafkaHandler
            void handleLMSEvent(LMSEvent event) {
                 ....
            }
    
            @KafkaHandler
            void handleYetAnotherEvent(YetAnotherEvent yetAnotherEvent) {
                ...
            }
    
            @KafkaHandler(isDefault = true)
            void handleDefault(@Payload Object unknown,
                                @Header(KafkaHeaders.OFFSET) long offset,
                                @Header(KafkaHeaders.RECEIVED_PARTITION) int partitionId,
                                @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
                logger.info("Server received unknown message {},{},{}", offset, partitionId, topic);
            }
        }
    

    Full code

    【讨论】:

      猜你喜欢
      • 2020-06-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-07-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多