Чтение Avro GenericRecord во FlinkJAVA

Программисты JAVA общаются здесь
Anonymous
Чтение Avro GenericRecord во Flink

Сообщение Anonymous »

Я хотел бы прочитать записи Avro (в моем задании Flink) из источника файла, без предоставления схемы (потому что меня просто интересуют пары ключ/значение).
Этот комментарий вверху AvroInputFormat в пакете org.apache.flink.formats.avro — именно то, что мне нужно:

Код: Выделить всё

/**
* Provides a {@link FileInputFormat} for Avro records.
*
* @param  the type of the result Avro record. If you specify {@link GenericRecord} then the
*     result will be returned as a {@link GenericRecord}, so you do not have to know the schema
*     ahead of time.
*/
Вот как я его вызываю (я пробовал и Scala, и Java):

Код: Выделить всё

val avroInputFormat = new AvroInputFormat(new Path("s3://.../records.avro"), classOf[GenericRecord])
env.createInput(avroInputFormat)
...
К сожалению, все попытки прочитать записи Avro с использованием этого FileInputFormat завершились с ошибкой UnsupportedOperationException. Поиск на различных сайтах в Интернете привел меня к выводу, что некоторые элементы GenericRecord не подлежат сериализации.
Кто-нибудь добился успеха в получении GenericRecord из Avro файл, без доступа к схеме? Если это актуально, я читаю не тему Kafka, а корзину S3 (так что это источник файла).
Вот часть трассировки исключения:

Код: Выделить всё

Caused by: com.esotericsoftware.kryo.KryoException: java.lang.UnsupportedOperationException
Serialization trace:
reserved (org.apache.avro.Schema$Field)
fieldMap (org.apache.avro.Schema$RecordSchema)
schema (org.apache.avro.generic.GenericData$Record)
at com.esotericsoftware.kryo.serializers.ObjectField.read(ObjectField.java:125)
at com.esotericsoftware.kryo.serializers.FieldSerializer.read(FieldSerializer.java:528)
at com.esotericsoftware.kryo.Kryo.readClassAndObject(Kryo.java:761)
at com.esotericsoftware.kryo.serializers.MapSerializer.read(MapSerializer.java:143)
at com.esotericsoftware.kryo.serializers.MapSerializer.read(MapSerializer.java:21)
Спасибо за любую помощь, которую вы можете предложить, чтобы помочь мне найти решение.
Если моя интерпретация комментария неверна (т. е. это не способ сделать то, чего я пытаюсь достичь), пожалуйста, дайте мне знать.

Подробнее здесь: https://stackoverflow.com/questions/785 ... d-in-flink

Вернуться в «JAVA»