如何在 Flink 中跳过损坏的消息?

mzl*_*zlo 0 apache-flink

如何在 Flink 中跳过损坏的消息

我有 DAG:KafkaSrcConsumer > FlatMap > Window > SinkFunction

现在,如果我在操作员“KafkaSrcConsumer”中收到来自 Kafka 的corruptedMessage,我想抛出/跳过该消息,并且我不想将该损坏的消息转发给下一个操作员“FlatMap”

我们如何在 Apache Flink 中实现这一点?

(注意:从 KafkaSrcConsumer 抛出异常将重新启动 flink 作业,我想避免这种情况,因为我只想跳过消息并移至下一条消息)

Dav*_*son 5

如果该deserialize(...)方法返回 null,则 Flink Kafka 消费者将默默地跳过损坏的消息。这在文档中有所描述。