如何在 Flink 中跳过损坏的消息
我有 DAG:KafkaSrcConsumer > FlatMap > Window > SinkFunction
现在,如果我在操作员“KafkaSrcConsumer”中收到来自 Kafka 的corruptedMessage,我想抛出/跳过该消息,并且我不想将该损坏的消息转发给下一个操作员“FlatMap”
我们如何在 Apache Flink 中实现这一点?
(注意:从 KafkaSrcConsumer 抛出异常将重新启动 flink 作业,我想避免这种情况,因为我只想跳过消息并移至下一条消息)
| 归档时间: |
|
| 查看次数: |
448 次 |
| 最近记录: |