Apache Flink读取Kafka的confluent-avro数据时如何跳过损坏记录
解决方案
Flink 的 avro-confluent 格式提供了和JSON格式完全对等的忽略解析错误配置,你只需要在建表的WITH参数中新增对应配置项即可:'avro-confluent.ignore-parse-errors' = 'true'
该配置开启后,Flink消费Kafka数据时遇到无法解析的损坏Avro记录会直接跳过,不会抛出异常中断任务运行,适配你无法控制上游写入数据的场景。
修改后的完整建表示例
CREATE TABLE recommendations ( ... ) WITH ( 'connector' = 'kafka', 'topic' = 'my_kafka_topic', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'testGroup', 'properties.security.protocol' = 'SASL_PLAINTEXT', 'properties.sasl.kerberos.service.name' = 'kafka', 'scan.startup.mode' = 'latest-offset', 'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'http://schema-registry-address', 'value.fields-include' = 'EXCEPT_KEY', -- 新增忽略解析错误配置 'value.avro-confluent.ignore-parse-errors' = 'true' );
注意事项
- 该配置从Flink 1.12版本开始正式支持,如果你使用的是更早的版本,需要自定义DeserializationSchema实现坏数据跳过逻辑
- 开启配置后默认不会留存坏数据,如果需要统计损坏记录数量、分析坏数据内容,可搭配Flink Table API的脏数据侧输出功能实现
内容的提问来源于stack exchange,提问作者Nitin Pandey
相关产品推荐
相关产品推荐

