Flink写入Kafka Avro Sink触发IllegalAccessException的解决方案问询
我用Flink Streaming从Kafka源主题读取JSON格式事件,去重后写入另一个Avro格式的Kafka主题,流程如下:
Kafka Topic(json格式) -> Flink Streaming(去重) -> Scala case class对象 -> Kafka Topic(Avro格式)
核心处理代码:
val sink = sinkProvider.getKafkaSink(brokerURL, targetTopic,kafkaTransactionMaxTimeoutMs, kafkaTransactionTimeoutMs) messageStream .map { record => convertJsonToExample(record) } .sinkTo(sink) .name("Example Kafka Avro Sink") .uid("Example-Kafka-Avro-Sink")
执行步骤
- 创建输出对应的Avro Schema
{ "type":"record", "name":"Example", "namespace":"ca.ix.dcn.test", "fields":[ { "name":"x", "type":"string" }, { "name":"y", "type":"long" } ] }
- 使用avro-hugger工具(版本1.2.1)基于Avro Schema生成对应SpecificRecord的Scala case class
- 使用Flink的
AvroSerializationSchema.forSpecific[Example]构建Kafka Sink,代码如下:
def getKafkaSink(brokers: String, targetTopic: String,transactionMaxTimeoutMs:String,transactionTimeoutMs:String) = { val schema = ReflectData.get.getSchema(classOf[Example]) val sink = KafkaSink.builder() .setBootstrapServers(brokers) .setProperty("transaction.max.timeout.ms",transactionMaxTimeoutMs) .setProperty("transaction.timeout.ms",transactionTimeoutMs) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(targetTopic) .setValueSerializationSchema(AvroSerializationSchema.forSpecific[Example](classOf[Example])) .setPartitioner(new FlinkFixedPartitioner()) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .build() sink }
运行异常
Caused by: org.apache.avro.AvroRuntimeException: java.lang.IllegalAccessException: Class org.apache.avro.specific.SpecificData can not access a member of class ca.ix.dcn.test with modifiers "private final" at org.apache.avro.specific.SpecificData.createSchema(SpecificData.java:405) at org.apache.avro.reflect.ReflectData.createSchema(ReflectData.java:734)
已知Flink存在相关BUG,未找到官方解决方案,现寻求问题的临时解决办法,以及使用AvroSerializationSchema(Specific/Generic模式)实现Flink Streaming Avro Sink的详细示例。
问题根源
异常是因为avro-hugger生成的Scala case class字段为private final,Avro的SpecificData类无法访问这些私有字段导致的。以下是几种可行的临时解决办法:
解决办法1:修改avro-hugger生成策略,生成公共字段的case class
在avro-hugger的配置中指定生成公共字段,以sbt插件为例:
avrohuggerSettings ++ Seq( avroSourceDirectories in Compile := Seq(new File("src/main/avro")), avroScalaCustomTypes := Map("string" -> "String", "long" -> "Long"), avroScalaGenerateCaseClassesWithPublicFields := true // 关键配置:生成公共访问权限的字段 )
重新生成Example case class后,字段权限变为公共,Avro即可正常访问。
解决办法2:改用Avro反射序列化(Reflect模式)
如果不想修改生成的case class,直接用AvroSerializationSchema.forReflect替代forSpecific,通过反射机制处理私有字段:
def getKafkaSink(brokers: String, targetTopic: String,transactionMaxTimeoutMs:String,transactionTimeoutMs:String) = { // 加载Avro Schema文件 val schema = new Schema.Parser().parse(new File("src/main/avro/example.avsc")) val sink = KafkaSink.builder() .setBootstrapServers(brokers) .setProperty("transaction.max.timeout.ms",transactionMaxTimeoutMs) .setProperty("transaction.timeout.ms",transactionTimeoutMs) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(targetTopic) // 改用反射序列化 .setValueSerializationSchema(AvroSerializationSchema.forReflect(classOf[Example], schema)) .setPartitioner(new FlinkFixedPartitioner()) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .build() sink }
解决办法3:手动编写符合Avro SpecificRecord要求的case class
放弃使用avro-hugger,手动实现继承SpecificRecordBase的case class,确保字段为公共权限并实现Avro要求的方法:
package ca.ix.dcn.test import org.apache.avro.Schema import org.apache.avro.specific.SpecificRecordBase case class Example(var x: String, var y: Long) extends SpecificRecordBase { // 必须提供无参构造函数 def this() = this("", 0L) override def getSchema: Schema = Example.SCHEMA$ override def get(field: Int): Any = field match { case 0 => x case 1 => y case _ => throw new IndexOutOfBoundsException() } override def put(field: Int, value: Any): Unit = field match { case 0 => x = value.asInstanceOf[String] case 1 => y = value.asInstanceOf[Long] case _ => throw new IndexOutOfBoundsException() } } object Example { val SCHEMA$: Schema = new Schema.Parser().parse("""{ "type":"record", "name":"Example", "namespace":"ca.ix.dcn.test", "fields":[ { "name":"x", "type":"string" }, { "name":"y", "type":"long" } ] }""") }
AvroSerializationSchema详细示例
1. Specific模式示例(正确版本)
使用avro-hugger生成带公共字段的case class后,Sink代码如下:
import org.apache.flink.connector.kafka.sink.KafkaSink import org.apache.flink.formats.avro.AvroSerializationSchema import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema import org.apache.flink.connector.kafka.sink.DeliveryGuarantee def getSpecificAvroKafkaSink(brokers: String, targetTopic: String): KafkaSink[Example] = { KafkaSink.builder[Example]() .setBootstrapServers(brokers) .setRecordSerializer( KafkaRecordSerializationSchema.builder[Example]() .setTopic(targetTopic) .setValueSerializationSchema(AvroSerializationSchema.forSpecific(classOf[Example])) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .build() } // 使用示例 messageStream .map(record => convertJsonToExample(record)) .sinkTo(getSpecificAvroKafkaSink("localhost:9092", "avro-output-topic")) .name("Specific-Avro-Kafka-Sink") .uid("specific-avro-kafka-sink-001")
2. Generic模式示例
无需生成case class,直接用GenericRecord处理:
import org.apache.avro.Schema import org.apache.avro.generic.GenericData import org.apache.avro.generic.GenericRecord import org.apache.flink.connector.kafka.sink.KafkaSink import org.apache.flink.formats.avro.AvroSerializationSchema import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema import com.fasterxml.jackson.databind.ObjectMapper // 加载Avro Schema val exampleSchema = new Schema.Parser().parse(new File("src/main/avro/example.avsc")) val objectMapper = new ObjectMapper() // JSON转GenericRecord def convertJsonToGenericRecord(jsonStr: String): GenericRecord = { val jsonNode = objectMapper.readTree(jsonStr) val record = new GenericData.Record(exampleSchema) record.put("x", jsonNode.get("x").asText()) record.put("y", jsonNode.get("y").asLong()) record } // 构建Generic模式Kafka Sink def getGenericAvroKafkaSink(brokers: String, targetTopic: String, schema: Schema): KafkaSink[GenericRecord] = { KafkaSink.builder[GenericRecord]() .setBootstrapServers(brokers) .setRecordSerializer( KafkaRecordSerializationSchema.builder[GenericRecord]() .setTopic(targetTopic) .setValueSerializationSchema(AvroSerializationSchema.forGeneric(schema)) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .build() } // 使用示例 messageStream .map(record => convertJsonToGenericRecord(record.value())) .sinkTo(getGenericAvroKafkaSink("localhost:9092", "avro-output-topic", exampleSchema)) .name("Generic-Avro-Kafka-Sink") .uid("generic-avro-kafka-sink-001")
内容的提问来源于stack exchange,提问作者S Mishra

