Flink写入AvroParquet Sink报Invalid lambda deserialization异常排查
故障根因
异常抛出在Parquet写入环节,Kafka读取正常说明Source逻辑无问题,核心触发点有三个:
AvroParquetWriters.forReflectRecord默认基于反射动态生成序列化逻辑,Scala编译生成的类结构、合成lambda方法无法被Java序列化机制正确识别,直接触发Invalid lambda deserialization错误MyEvent类不符合Avro反射序列化要求:未给java.util.Date类型指定序列化编码,Scala默认生成的字段访问器无法被Avro ReflectData正确解析- KafkaSource配置存在冲突:同时设置了批模式的
setBounded(OffsetsInitializer.latest)和流模式起始偏移量配置,后续运行会出现消费终止、偏移量重置等问题
修复方案
1. 改造MyEvent类适配Avro序列化
给日期字段加Avro编码注解,保证类结构符合POJO规范:
import com.fasterxml.jackson.annotation.JsonProperty import org.apache.avro.reflect.AvroEncode import org.apache.avro.reflect.DateAsLongEncoding import java.io.Serializable import java.util.Date class MyEvent() extends Serializable { @JsonProperty("id") var id: String = _ // 指定日期类型用long型时间戳序列化,避免Avro默认不支持Date类型的问题 @AvroEncode(using = classOf[DateAsLongEncoding]) @JsonProperty("timestamp") var timestamp: Date = _ }
2. 修正KafkaSource配置
流模式下移除setBounded配置,该配置仅用于批作业消费到指定偏移量后自动停止:
val builder = KafkaSource.builder[MyEvent] builder.setBootstrapServers("localhost:29092") builder.setProperty("partition.discovery.interval.ms", "10000") builder.setTopics("myevent") // 移除冲突的setBounded配置 builder.setStartingOffsets(OffsetsInitializer.earliest) builder.setDeserializer(KafkaRecordDeserializationSchema.of(new MyEventSchema)) val source = builder.build()
3. 重构StreamingFileSink构造逻辑
提前生成Avro Schema,用显式声明的Writer工厂构造Sink,避免运行时动态生成lambda无法序列化:
import org.apache.flink.formats.parquet.avro.AvroParquetWriterFactory import org.apache.avro.reflect.ReflectData // 提前加载MyEvent对应的Avro Schema val myEventSchema = ReflectData.get().getSchema(classOf[MyEvent]) // 显式构造Parquet Writer工厂,绑定Schema val parquetWriterFactory = AvroParquetWriterFactory .forReflectRecord[MyEvent](classOf[MyEvent]) .withSchema(myEventSchema) .build() val sink: StreamingFileSink[MyEvent] = StreamingFileSink .forBulkFormat(new Path("hdfs://localhost:50070/mydata"), parquetWriterFactory) .build()
如果仍有序列化报错,在执行环境中注册MyEvent的序列化器:
val env = StreamExecutionEnvironment.getExecutionEnvironment env.getConfig.registerTypeWithKryoSerializer( classOf[MyEvent], classOf[com.esotericsoftware.kryo.serializers.JavaSerializer] )
4. 环境适配检查
- 进入namenode容器提前创建HDFS写入路径并开放权限:
hdfs dfs -mkdir -p /mydata hdfs dfs -chmod 777 /mydata - 保证程序依赖的Avro、Parquet、Hadoop客户端版本和docker部署的组件版本一致,避免版本冲突导致的写入失败
- 如果最终需要对接Hive查询,直接使用Flink Hive Connector写入Hive表,无需手动写Parquet文件再关联元数据,稳定性更高
内容的提问来源于stack exchange,提问作者Sjoerd222888
相关产品推荐
相关产品推荐

