如何在Flink Producer中捕获FlinkRuntimeException跳过失败消息防止作业重启
是否存在可行方案能够捕获FlinkRuntimeException并对其进行妥善处理,避免Flink作业触发重启?
抛出FlinkRuntimeException的代码示例
int numberOfRows = 1; int rowsPerSecond = 1; DataStream<String> stream = environment.addSource( new DataGeneratorSource<>( RandomGenerator.stringGenerator(1050000), // max.message.bytes=1048588 rowsPerSecond, (long) numberOfRows), TypeInformation.of(String.class)) .setParallelism(1) .name("string-generator"); KafkaSinkBuilder<String> builder = KafkaSink.<String>builder() .setBootstrapServers("localhost:9092") .setDeliverGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) .setRecordSerializer( KafkaRecordSerializationSchema.builder().setTopic("test.output") .setValueSerializationSchema(new SimpleStringSchema()) .build()); KafkaSink<String> sink = builder.build(); stream.sinkTo(sink).setParallelism(1).name("output-producer");
异常栈追踪信息
2022-06-02/14:01:45.066/PDT [flink-akka.actor.default-dispatcher-4] INFO output-producer: Writer -> output-producer: Committer (1/1) (a66beca5a05c1c27691f7b94ca6ac025)
switched from RUNNING to FAILED on 271b1b90-7d6b-4a34-8116-3de6faa8a9bf @ 127.0.0.1 (dataPort=-1).
org.apache.flink.util.FlinkRuntimeException: Failed to send data to Kafka null with FlinkKafkaInternalProducer{transactionalId='null', inTransaction=false, closed=false}
at org.apache.flink.connector.kafka.sink.KafkaWriter$WriterCallback.throwException(KafkaWriter.java:440) ~[flink-connector-kafka-1.15.0.jar:1.15.0]
at org.apache.flink.connector.kafka.sink.KafkaWriter$WriterCallback.lambda$onCompletion$0(KafkaWriter.java:421) ~[flink-connector-kafka-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:90) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMailsNonBlocking(MailboxProcessor.java:353) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:317) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:201) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:804) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:753) ~[flink-streaming-java-1.15.0.jar:1.15.0]
at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:948) ~[flink-runtime-1.15.0.jar:1.15.0]
at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:927) ~[flink-runtime-1.15.0.jar:1.15.0]
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:741) ~[flink-runtime-1.15.0.jar:1.15.0]
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:563) ~[flink-runtime-1.15.0.jar:1.15.0]
at java.lang.Thread.run(Thread.java:748) ~[?:1.8.0_292]
Caused by: org.apache.kafka.common.errors.RecordTooLargeException: The message is 1050088 bytes when serialized which is larger than 1048576, which is the value of the max.request.size configuration.
首先明确:这个异常是Kafka Producer异步发送的回调里抛出的,被Kafka Sink包装成FlinkRuntimeException后提交到Task的Mailbox线程触发失败,普通业务算子外层的try-catch根本抓不到,别在map/process之类的算子外面包try-catch白费功夫。
可行的处理方案有三个,按推荐优先级排序:
- 方案1:调整Kafka相关配置,从根源消除异常(成本最低,无数据风险)
报错信息已经明确给出根因:单条消息序列化后大小为1050088字节,超过了Kafka Producer默认的max.request.size1MB阈值。直接在构造KafkaSink时传入Producer配置调大阈值即可,注意要同步保证Kafka Broker端的message.max.bytes、对应Topic的max.message.bytes配置不小于这个值,否则还是会报错。配置示例:// 把单条请求最大大小调整为2MB,可根据实际消息大小按需设置 builder.setProperty("max.request.size", "2097152"); - 方案2:在Sink前做前置校验,拦截非法超大消息
如果业务规范不允许调大Kafka消息阈值,就在写入Kafka前加一个处理算子,提前计算消息序列化后的大小:- 超过阈值的消息直接打错误日志记录内容后丢弃
- 或者把超限消息路由到侧输出流,单独写入HDFS、对象存储等支持大消息的存储做兜底,不进入Kafka链路
- 也可以对超限消息做拆分、压缩,保证进入KafkaSink的消息全部符合大小要求,从源头避免异常触发。
- 方案3:自定义Kafka发送逻辑吞掉异常(不推荐,有丢数风险)
如果业务允许部分消息发送失败,可以自定义Kafka Producer Sink,重写异步发送的回调逻辑,在回调里捕获RecordTooLargeException这类你认为可容忍的异常,只打日志不向上抛出,作业就不会触发重启。注意这种方式会直接丢失发送失败的消息,没有兜底逻辑的话绝对不要用。
补充:Flink作业重启的触发逻辑是只要算子链中抛出未捕获的异常向上传到Task层,就会触发失败重启,所以处理这类问题的核心要么是让异常根本不产生,要么是在异常抛到Task层之前拦住它。
内容的提问来源于stack exchange,提问作者Dtmy

