【问题标题】:How to deserialize avro data using Apache Beam (KafkaIO)如何使用 Apache Beam (KafkaIO) 反序列化 avro 数据
【发布时间】:2019-09-13 08:14:04
【问题描述】:

我只看到一个帖子包含有关我提到的主题的信息,即: How to Deserialising Kafka AVRO messages using Apache Beam

但是,在尝试了一些 kafkaserializers 变体之后,我仍然无法反序列化 kafka 消息。这是我的代码:

public class Readkafka {
    private static final Logger LOG = LoggerFactory.getLogger(Readkafka.class);

    public static void main(String[] args) throws IOException {
        // Create the Pipeline object with the options we defined above.
        Pipeline p = Pipeline.create(
                PipelineOptionsFactory.fromArgs(args).withValidation().create());
       PTransform<PBegin, PCollection<KV<action_states_pkey, String>>> kafka =
                KafkaIO.<action_states_pkey, String>read()
                    .withBootstrapServers("mybootstrapserver")
                    .withTopic("action_States")
                    .withKeyDeserializer(MyClassKafkaAvroDeserializer.class)
                    .withValueDeserializer(StringDeserializer.class)
                    .updateConsumerProperties(ImmutableMap.of("schema.registry.url", (Object)"schemaregistryurl"))
                    .withMaxNumRecords(5)
                    .withoutMetadata();


        p.apply(kafka)
            .apply(Keys.<action_states_pkey>create())
}

MyClassKafkaAvroDeserilizer 在哪里

public class MyClassKafkaAvroDeserializer extends
AbstractKafkaAvroDeserializer implements Deserializer<action_states_pkey> {

@Override
public void configure(Map<String, ?> configs, boolean isKey) {
    configure(new KafkaAvroDeserializerConfig(configs));
}

@Override
public action_states_pkey deserialize(String s, byte[] bytes) {
    return (action_states_pkey) this.deserialize(bytes);
}

@Override
public void close() {} }

action_states_pkey 类是使用 avro 工具生成的代码

java -jar pathtoavrotools/avro-tools-1.8.1.jar compile schema pathtoschema/action_states_pkey.avsc destination path

action_states_pkey.avsc 的字面意思

{"type":"record","name":"action_states_pkey","namespace":"namespace","fields":[{"name":"ad_id","type":["null","int"]},{"name":"action_id","type":["null","int"]},{"name":"state_id","type":["null","int"]}]}

使用此代码,我得到了错误:

Caused by: java.lang.ClassCastException: org.apache.avro.generic.GenericData$Record cannot be cast to my.mudah.beam.test.action_states_pkey
    at my.mudah.beam.test.MyClassKafkaAvroDeserializer.deserialize(MyClassKafkaAvroDeserializer.java:20)
    at my.mudah.beam.test.MyClassKafkaAvroDeserializer.deserialize(MyClassKafkaAvroDeserializer.java:1)
    at org.apache.beam.sdk.io.kafka.KafkaUnboundedReader.advance(KafkaUnboundedReader.java:221)
    at org.apache.beam.sdk.io.BoundedReadFromUnboundedSource$UnboundedToBoundedSourceAdapter$Reader.advanceWithBackoff(BoundedReadFromUnboundedSource.java:279)
    at org.apache.beam.sdk.io.BoundedReadFromUnboundedSource$UnboundedToBoundedSourceAdapter$Reader.start(BoundedReadFromUnboundedSource.java:256)
    at com.google.cloud.dataflow.worker.WorkerCustomSources$BoundedReaderIterator.start(WorkerCustomSources.java:592)
    ... 14 more

尝试将 Avro 数据映射到我的自定义类时似乎出错了?

或者,我尝试了以下代码:

        PTransform<PBegin, PCollection<KV<action_states_pkey, String>>> kafka =
                KafkaIO.<action_states_pkey, String>read()
                    .withBootstrapServers("bootstrapserver")
                    .withTopic("action_states")
                    .withKeyDeserializerAndCoder((Class)KafkaAvroDeserializer.class, AvroCoder.of(action_states_pkey.class))
                    .withValueDeserializer(StringDeserializer.class)
                    .updateConsumerProperties(ImmutableMap.of("schema.registry.url", (Object)"schemaregistry"))
                    .withMaxNumRecords(5)
                    .withoutMetadata();


        p.apply(kafka);
            .apply(Keys.<action_states_pkey>create())
//            .apply("ExtractWords", ParDo.of(new DoFn<action_states_pkey, String>() {
//                @ProcessElement
//                public void processElement(ProcessContext c) {
//                  action_states_pkey key = c.element();
//                    c.output(key.getAdId().toString());
//                }
//            }));

在我尝试打印数据之前不会给我任何错误。我必须验证我是否以一种或另一种方式成功读取数据,所以我的意图是在控制台中记录数据。如果我取消注释注释部分,我会再次收到相同的错误:

SEVERE: 2019-09-13T07:53:56.168Z: java.lang.ClassCastException: org.apache.avro.generic.GenericData$Record cannot be cast to my.mudah.beam.test.action_states_pkey
    at my.mudah.beam.test.Readkafka$1.processElement(Readkafka.java:151)

另外需要注意的是,如果我指定:

.updateConsumerProperties(ImmutableMap.of("specific.avro.reader", (Object)"true"))

总是给我一个错误

Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id 443
Caused by: org.apache.kafka.common.errors.SerializationException: Could not find class NAMESPACE.action_states_pkey specified in writer's schema whilst finding reader's schema for a SpecificRecord.

我的方法似乎有问题? 如果有人有使用 Apache Beam 从 Kafka Streams 读取 AVRO 数据的经验,请帮帮我。非常感谢。

这是我的包的快照,其中还包含架构和类: package/working path details

谢谢。

【问题讨论】:

  • 您似乎在使用 Confluent Schema Registry?那么为什么不使用他们现有的 Avro 解串器呢?如果你这样做了,那么这与你拥有的其他反序列化器实现不兼容......换句话说,数据实际上是如何产生的?使用模式注册表或原始 Avro 数据?
  • 我假设我已经在使用 confluent 的反序列化器了?导入 io.confluent.kafka.serializers.KafkaAvroDeserializer;这是我正在使用的图书馆
  • 抱歉,我不太明白您所说的“数据实际上是如何产生的?”使用模式注册表或原始 Avro 数据? ' 我从架构注册表中获取架构,将其编译成 java.class,然后继续
  • 抱歉,我以为您编写了一个不同的反序列化程序,最初并没有使用注册表。我认为你的第二种方法是正确的。我不确定为什么命名空间被大写
  • 当然。出于隐私原因试图掩盖它。这是错误: 原因:org.apache.kafka.common.errors.SerializationException:反序列化 id 443 的 Avro 消息时出错 原因:org.apache.kafka.common.errors.SerializationException:找不到类 com.dattran.bottledwater .dbschema.public.action_states_pkey 在作者模式中指定,同时为特定记录查找读者模式。我猜 com.dattran.bottledwater.dbschema.public 是架构在架构注册表中的位置?

标签: java apache-kafka avro apache-beam confluent-schema-registry


【解决方案1】:

公共类 MyClassKafkaAvroDeserializer 扩展 AbstractKafkaAvroDeserializer

您的课程正在扩展AbstractKafkaAvroDeserializer,它返回GenericRecord

你需要convert the GenericRecord to your custom object

如以下答案之一所述,为此使用SpecificRecord

/**
 * Extends deserializer to support ReflectData.
 *
 * @param <V>
 *     value type
 */
public abstract class ReflectKafkaAvroDeserializer<V> extends KafkaAvroDeserializer {

  private Schema readerSchema;
  private DecoderFactory decoderFactory = DecoderFactory.get();

  protected ReflectKafkaAvroDeserializer(Class<V> type) {
    readerSchema = ReflectData.get().getSchema(type);
  }

  @Override
  protected Object deserialize(
      boolean includeSchemaAndVersion,
      String topic,
      Boolean isKey,
      byte[] payload,
      Schema readerSchemaIgnored) throws SerializationException {

    if (payload == null) {
      return null;
    }

    int schemaId = -1;
    try {
      ByteBuffer buffer = ByteBuffer.wrap(payload);
      if (buffer.get() != MAGIC_BYTE) {
        throw new SerializationException("Unknown magic byte!");
      }

      schemaId = buffer.getInt();
      Schema writerSchema = schemaRegistry.getByID(schemaId);

      int start = buffer.position() + buffer.arrayOffset();
      int length = buffer.limit() - 1 - idSize;
      DatumReader<Object> reader = new ReflectDatumReader(writerSchema, readerSchema);
      BinaryDecoder decoder = decoderFactory.binaryDecoder(buffer.array(), start, length, null);
      return reader.read(null, decoder);
    } catch (IOException e) {
      throw new SerializationException("Error deserializing Avro message for id " + schemaId, e);
    } catch (RestClientException e) {
      throw new SerializationException("Error retrieving Avro schema for id " + schemaId, e);
    }
  }
}

以上内容抄自https://stackoverflow.com/a/39617120/2534090

https://stackoverflow.com/a/42514352/2534090

【讨论】:

  • 感谢您的洞察,如果我已经有了架构 (.avsc),代码会是什么样子?您编写的代码是如果我没有正确的架构?它从 avro 数据中获取模式?
  • 我还认为添加 .updateConsumerProperties(ImmutableMap.of("specific.avro.reader", (Object)"true")) 使其能够读取为特定记录?还是我弄错了?但我最终得到错误找不到作者架构中指定的类 NAMESPACE.action_states_pkey,同时为特定记录找到读者架构。
  • @NicholasLeongZhiHao 如果你有 .avsc 文件,那么你可以使用 apache avro zip 文件附带的 avro-tools 从中生成相应的 POJO 文件。这些 POJO 是我们包含在 SpecificRecord 中的那些。它们包含一些额外的字段,看起来不像我们通常编写的典型 POJO。
  • 感谢您的回复。如前所述,我已经使用 avro 工具生成了 POJO。这里的名称为 action_states_pkey.java。关于我应该如何处理 POJO 和您建议的代码的任何指示?非常感谢
  • @NicholasLeongZhiHao 你得到的错误是什么?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-07-12
  • 1970-01-01
  • 2023-03-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多