Этот комментарий вверху 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.
*/
Код: Выделить всё
val avroInputFormat = new AvroInputFormat(new Path("s3://.../records.avro"), classOf[GenericRecord])
env.createInput(avroInputFormat)
...
Кто-нибудь добился успеха в получении 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