You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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")

执行步骤

  1. 创建输出对应的Avro Schema
{
   "type":"record",
   "name":"Example",
   "namespace":"ca.ix.dcn.test",
   "fields":[
      {
         "name":"x",
         "type":"string"
      },
      {
         "name":"y",
         "type":"long"
      }
   ]
}
  1. 使用avro-hugger工具(版本1.2.1)基于Avro Schema生成对应SpecificRecord的Scala case class
  2. 使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 05:01:11