【问题标题】:What is better way put several Event Types in the same Kafka topic?将多个事件类型放在同一个 Kafka 主题中的更好方法是什么?
【发布时间】:2020-01-02 15:48:24
【问题描述】:

假设有两种类型 T1 & T2 和一个主题 TT1T2 都必须进入主题 T(出于某种原因)。有什么方法可以实现这一目标?哪个更好?

其中一种方法是利用继承,我们可以定义一个基类,然后子类可以扩展它。在我们的例子中,我们可以定义一个基类 TB,然后 T1 & T2 可以扩展 TB

基类 (TB)

package poc.kafka.domain;

    import java.io.Externalizable;
    import java.io.IOException;
    import java.io.ObjectInput;
    import java.io.ObjectOutput;

    import lombok.AllArgsConstructor;
    import lombok.NoArgsConstructor;
    import lombok.ToString;
    import lombok.extern.java.Log;

    @ToString
    @AllArgsConstructor
    @NoArgsConstructor
    @Log
    public class Animal implements Externalizable {
        public String name;

        public void whoAmI() {
            log.info("I am an Animal");
        }

        @Override
        public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException {
            name = (String) in.readObject();
        }

        @Override
        public void writeExternal(ObjectOutput out) throws IOException {
            out.writeObject(name);
        }
    }

派生类 (T1)

package poc.kafka.domain;

import java.io.Externalizable;
import java.io.IOException;
import java.io.ObjectInput;
import java.io.ObjectOutput;

import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import lombok.ToString;
import lombok.extern.java.Log;

@Log
@Setter
@Getter
@AllArgsConstructor
@NoArgsConstructor
@ToString
public class Cat extends Animal implements Externalizable {
    private int legs;

    public void whoAmI() {
        log.info("I am a Cat");
    }

    @Override
    public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException {
        super.readExternal(in);
        legs = in.readInt();
    }

    @Override
    public void writeExternal(ObjectOutput out) throws IOException {
        super.writeExternal(out);
        out.writeInt(legs);
    }
}

派生类 (T2)

package poc.kafka.domain;

import java.io.Externalizable;
import java.io.IOException;
import java.io.ObjectInput;
import java.io.ObjectOutput;

import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import lombok.ToString;
import lombok.extern.java.Log;

@Log
@Setter
@Getter
@AllArgsConstructor
@NoArgsConstructor
@ToString
public class Dog extends Animal implements Externalizable {
    private int legs;

    public void whoAmI() {
        log.info("I am a Dog");
    }

    @Override
    public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException {
        super.readExternal(in);
        legs = in.readInt();
    }

    @Override
    public void writeExternal(ObjectOutput out) throws IOException {
        super.writeExternal(out);
        out.writeInt(legs);
    }
}

反序列化器

package poc.kafka.domain.serialization;

import org.apache.commons.lang3.SerializationUtils;
import org.apache.kafka.common.serialization.Deserializer;

import poc.kafka.domain.Animal;

public class AnimalDeserializer implements Deserializer<Animal> {

    @Override
    public Animal deserialize(String topic, byte[] data) {
        return SerializationUtils.deserialize(data);
    }

}

序列化器

package poc.kafka.domain.serialization;

import org.apache.commons.lang3.SerializationUtils;
import org.apache.kafka.common.serialization.Serializer;

import poc.kafka.domain.Animal;

public class AnimalSerializer implements Serializer<Animal> {

    @Override
    public byte[] serialize(String topic, Animal data) {
        return SerializationUtils.serialize(data);
    }

}

然后我们可以发送 T1 & T2 如下所示

IntStream.iterate(0, i -> i + 1).limit(10).forEach(i -> {
            if (i % 2 == 0)
                producer.send(new ProducerRecord<Integer, Animal>("T", i, new Dog(i)));
            else
                producer.send(new ProducerRecord<Integer, Animal>("gs3", i, new Cat(i)));
        });

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    最好的方法是创建自定义分区。

    通过partitionKey将每条消息生成到不同的分区

    这是default implementation,你需要实现你的分区逻辑。

    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
            if (keyBytes == null) {
                return stickyPartitionCache.partition(topic, cluster);
            } 
            List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
            int numPartitions = partitions.size();
            // hash the keyBytes to choose a partition
            return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
        }
    

    查看tutorial 了解更多示例。

    这是来自 kafka 的一段关于何时选择服装分区的权威指南。

    实现自定义分区策略 到目前为止,我们已经讨论了 默认分区器的特征,这是最常见的一种 用过的。但是,Kafka 并不仅限于散列分区,而且 有时有充分的理由对数据进行不同的分区。为了 例如,假设您是 B2B 供应商和最大的客户 是一家制造名为 Bananas 的手持设备的公司。 假设你与客户“Banana”做了很多生意,以至于 超过 10% 的日常交易是与该客户进行的。如果你使用 默认哈希分区,香蕉记录将被分配到 与其他帐户相同的分区,导致一个分区被 大约是其余部分的两倍。这可能会导致服务器用完 空间,减慢处理速度等。我们真正想要的是给予 Banana 自己的分区然后使用哈希分区来映射 其余帐户到分区。

    【讨论】:

    • 我认为这不是好方法。为什么他们应该去不同的分区?即使他们去不同的分区,同一个 Kafka Consumer 也可能会收到他们。
    • @wardziniak 消费者需要按业务逻辑用例从正确的分区中读取数据
    • 使用这种方法,您不会使用自动组管理器,并且您的应用程序不能容错
    • 这样你会失去很多kafkas的保证。
    【解决方案2】:

    最简单的方法是使用您的自定义org.apache.kafka.common.serialization.Serializer,它将能够处理这两种类型的事件。两种类型的事件都应该继承自相同的类型/基类。

    示例代码如下所示:

    public class CustomSerializer implements Serializer<T> {
    
        public void configure(Map<String, ?> configs, boolean isKey) {
            // nothing to do
        }
    
        public byte[] serialize(String topic, T data) {
            // serialization
            return null;
        }
    
        public void close() {
            // nothing to do
        }
    } 
    

    【讨论】:

    • 这是我选择的。继承是这里的棘手部分。如果我要重构我的应用程序,我会从 @JsonSubType 转移到 Mixin-Annotations。
    【解决方案3】:

    这可能不是问题的直接答案,而是提议重新考虑这里的某些方面,这可能会解决最初的问题。

    首先,尽管 Kafka 能够支持任何数据格式,但对于可序列化的二进制格式,我建议使用 Apache Avro,而不是序列化的 Java 对象。

    使用 Avro,您将获得紧凑的二进制、与语言无关的数据类型和广泛的工具集所带来的所有好处。例如,有 CLI 工具可以在 Avro 中读取带有内容的 Kafka 主题,但我不知道有哪一个工具能够反序列化 Java 对象。

    您可以阅读有关 Avro 本身的信息 here

    在这个 SO 问题 here

    中还可以找到关于为什么使用 Avro 的一些很好的见解

    第二。您的问题标题是关于 Event 类型的,但判断描述可能暗示“如何通过单个 Kafka 主题处理不同的 data 类型”。如果事件之间的区别只是事件类型——例如,Click、Submit、LogIn、LogOut 等等——那么你可以在里面保留一个这种类型的 enum 字段,否则使用通用容器对象。

    如果这些事件应携带的数据负载的结构存在差异,那么再次使用 Avro,您可以使用 Union 类型解决它。

    最后,如果数据差异如此之大,以至于这些事件基本上是不同的数据结构,没有什么共同点 - 请选择不同的 Kafka 主题

    尽管能够在同一个主题中使用不同的分区来发送不同的数据类型,但它确实只会在未来引起维护方面的麻烦,并且正如其他回复中正确指出的那样,对扩展的限制。因此,对于这种情况,如果可以选择处理不同的主题 - 最好这样做。

    【讨论】:

      【解决方案4】:

      如果没有继承的概念,比如数据不一样

      Animal -> Cat
      Animal -> Dog
      

      那么另一种方法是使用包装器。

      public class DataWrapper
      {
      private Object data;
      private EventType type;
             // getter and setters omitted for the sake of brevity
      }
      

      将所有事件放在包装器对象中,并用EventType 区分每个事件,例如enum

      然后您可以以正常方式对其进行序列化(正如您在问题中发布的那样),在反序列化时您可以检查EventType,然后根据EventType将其委托给相应的事件处理器

      此外,为了确保您的 DataWrapper 不会包装所有类型的数据,即应仅用于特定类型的数据,那么您可以使用Marker 接口,并使您将推送到主题的对象的所有类都实现此接口。

      例如,

      interface MyCategory {
      }
      

      然后你的自定义类可以有例如,

      class MyEvent implements MyCategory {
      }
      

      DataWrapper 中你可以拥有..

      public class DataWrapper<T extends MyCategory> {
      private T data;
      private EventType type;
                  // getters and setters omitted for the sake of brevity
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2018-12-28
        • 2018-11-02
        • 2021-03-27
        • 2021-06-08
        • 2019-09-26
        • 2023-03-07
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多