【发布时间】:2020-02-09 04:31:39
【问题描述】:
我正在尝试查找有关使用两种不同 Avro 类型发送 Kafka 消息的性能和(缺点)优势的一些信息。 根据我的研究,可以创建一个基于 avro 的 Kafka 消息的有效负载:
任一:
GenericRecord 可以通过调用 new GenericData.Record 并将从 Schema Registry 读取的模式作为参数传递来创建其实例:
大致:
private CachedSchemaRegistryClient schemaRegistryClient;
private Schema valueSchema;
// Read a schema
//…
this.valueSchema = schemaRegistryClient.getBySubjectAndID("TestTopic-value",1);
// Define a generic record according to the loaded schema
GenericData.Record record = new GenericData.Record(valueSchema);
// Send to kafka
ListenableFuture<SendResult<String, GenericRecord>> res;
res = avroKafkaTemplate
.send(MessageBuilder
.withPayload(record)
.setHeader(KafkaHeaders.TOPIC, TOPIC)
.setHeader(KafkaHeaders.MESSAGE_KEY, record.get("id"))
.build());
或:
扩展了 SpecificRecordBase 并在 Maven 的帮助下生成的类(来自包含 Avro 架构的文件)
/..
public class MyClass extends org.apache.avro.specific.SpecificRecordBase implements org.apache.avro.specific.SpecificRecord
/..
MyClass myAvroClass = new MyClass();
ListenableFuture<SendResult<String, MyClass>> res;
res = avroKafkaTemplate
.send(MessageBuilder
.withPayload(myAvroClass)
.setHeader(KafkaHeaders.TOPIC, TOPIC)
.setHeader(KafkaHeaders.MESSAGE_KEY, myAvroClass.getId())
.build());
当一段包含扩展GenericRecord的类的实例的代码被调试时,可以看到其中包含一个架构。
关于这个问题,我有几个问题:
如果我向 Kafka 发送 GenericRecord 实例,是否也发送了底层架构?
如果没有,什么时候下架?哪个类/方法负责从 GenericRecord 中提取字节并删除底层架构,使其不与有效负载一起发送? 如果是,那么架构注册表的意义何在?-
如果类扩展了 SpecificRecord,底层架构也会被发送,不是吗?这意味着,如果我使用一个接收 Kafka 消息并计算其字节数的函数,我应该期望特定记录消息中的字节数比通用记录消息中的字节多,对吧?
-
SpecificRecord 实例给了我更多的控制权,而且使用起来更不容易出错。如果模式不是与 GenericRecord 一起发送而与 SpecificRecord 一起发送,那么我们需要权衡。 一方面(SpecificRecord),使用简单,因为有清晰的 API 可用(不必记住所有字段,并编写 get("X")、get("Y") 等) ,另一方面,有效负载的大小会增加,因为模式必须与它一起发送。如果我有一个相对较大的架构(50 个字段),我应该选择在架构注册表的帮助下发送 GenericRecords,否则性能会受到负面影响,因为架构必须随每条消息一起发送,对吗?
【问题讨论】:
标签: java apache-kafka avro confluent-platform confluent-schema-registry