【问题标题】:how to send avro schema ONLY once in kafka如何在kafka中只发送一次avro模式
【发布时间】:2017-01-28 22:19:32
【问题描述】:

我正在使用以下代码(不是真的,但让我们假设它)创建一个模式并由生产者将其发送给 kafka。

public static final String USER_SCHEMA = "{"
        + "\"type\":\"record\","
        + "\"name\":\"myrecord\","
        + "\"fields\":["
        + "  { \"name\":\"str1\", \"type\":\"string\" },"
        + "  { \"name\":\"str2\", \"type\":\"string\" },"
        + "  { \"name\":\"int1\", \"type\":\"int\" }"
        + "]}";

public static void main(String[] args) throws InterruptedException {
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");

    Schema.Parser parser = new Schema.Parser();
    Schema schema = parser.parse(USER_SCHEMA);
    Injection<GenericRecord, byte[]> recordInjection = GenericAvroCodecs.toBinary(schema);

    KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props);

    for (int i = 0; i < 1000; i++) {
        GenericData.Record avroRecord = new GenericData.Record(schema);
        avroRecord.put("str1", "Str 1-" + i);
        avroRecord.put("str2", "Str 2-" + i);
        avroRecord.put("int1", i);

        byte[] bytes = recordInjection.apply(avroRecord);

        ProducerRecord<String, byte[]> record = new ProducerRecord<>("mytopic", bytes);
        producer.send(record);

        Thread.sleep(250);

    }

    producer.close();
}

问题是代码只允许我使用此架构发送 1 条消息。然后我需要更改架构名称以发送下一条消息......所以名称字符串现在是随机生成的,所以我可以发送更多消息。这是一个 hack,所以我想知道正确的方法。

我还研究了如何在没有架构的情况下发送消息(即,已经向 kafka 发送了 1 条带有架构的消息,现在所有其他消息都不再需要架构了)——但 new GenericData.Record(..) 需要一个架构参数。如果为 null 则会抛出错误。

那么将 avro 模式消息发送到 kafka 的正确方法是什么?

这是另一个代码示例 - 与我的完全相同:
https://github.com/confluentinc/examples/blob/kafka-0.10.0.1-cp-3.0.1/kafka-clients/producer/src/main/java/io/confluent/examples/producer/ProducerExample.java

它也没有显示如何在不设置架构的情况下发送。

【问题讨论】:

    标签: apache-kafka avro kafka-producer-api confluent-platform apache-kafka-connect


    【解决方案1】:

    我没看懂这行:

    问题是代码只允许我发送 1 条消息 架构。然后我需要更改架构名称才能发送 下一条消息。

    在您提供的这两个示例和您提供的融合示例中,架构都没有发送到 Kafka。

    在您提供的示例中,用于创建 GenericRecord 对象的架构。您提供架构,因为您想针对某个架构验证记录(例如,验证您只能将整数 int1 字段放入 GenericRecord 对象中)。

    在您的代码中,唯一的区别是您决定将数据序列化为 byte[],这可能不是必需的,因为您可以将此责任委托给 KafkaAvroSerializer,正如您在 confluent 示例中所见。

    GenericRecord 是一个 Avro 对象,它不是 Kafka 强制执行的。如果您想将任何类型的对象发送到 Kafka(带有或不带有模式),您只需要创建(或使用现有的)序列化程序,它将您的对象转换为 byte[] 并在您创建的属性中设置此序列化程序生产者。

    通常最好使用 Avro 消息本身发送指向架构的指针。您可以在以下链接中找到原因: http://www.confluent.io/blog/schema-registry-kafka-stream-processing-yes-virginia-you-really-need-one/

    【讨论】:

    • 谢谢,我会改写我的问题:我想从 kafka 写到 hdfs。从我所见,我需要以 avro 格式保存 kafka 中的数据,否则 hdfs 不会写任何东西。现在我正在尝试以 avro 格式将数据写入 kafka。我想注册一次模式(或根本不注册)并发送消息。你有这个的示例代码吗?在向 kafka 发送数据之前,是否需要先向模式服务器注册模式?我可以像上面的代码那样以 avro 格式将数据发送到 kafka,而无需在消息中添加模式。
    • 忘了提到每次我发送消息时都需要更改模式名称,否则我会收到错误消息:`Can't redefine: myrecord`。这是一个错误吗?
    • HDFS 只是一个文件系统。您可以选择将数据写入 HDFS 的内容和方式。您计划如何将数据从 Kafka 移动到 HDFS?你需要有一个负责它的流程。也许您计划使用的进程强制使用 Avro,但这不是 HDFS 的先决条件。您可以写入各种格式的数据。您可以在这里看到将字符串写入 Kafka 的示例:cwiki.apache.org/confluence/display/KAFKA/… 稍后您可以使用这些字符串并使用 HDFS API 将它们写入 HDFS。
    猜你喜欢
    • 2022-11-08
    • 2018-01-26
    • 2021-11-25
    • 2016-11-16
    • 1970-01-01
    • 1970-01-01
    • 2020-08-05
    • 2022-01-12
    • 1970-01-01
    相关资源
    最近更新 更多