【发布时间】:2018-12-13 22:33:16
【问题描述】:
我们的 Flink 流式工作流将消息发布到 Kafka。 KafkaProducer 的“重试”机制在将消息添加到其内部缓冲区之前不会启动。
如果在此之前有异常,KafkaProducer 将抛出该异常,并且似乎 Flink 没有处理该异常。在这种情况下,将会丢失数据。
相关Flink代码(FlinkKafkaProducerBase):
if (logFailuresOnly) {
callback = new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception e) {
if (e != null) {
LOG.error("Error while sending record to Kafka: " + e.getMessage(), e);
}
acknowledgeMessage();
}
};
}
else {
callback = new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception != null && asyncException == null) {
asyncException = exception;
}
acknowledgeMessage();
}
};
}
以下是我们已确定会导致数据丢失的场景:
-
所有的 kafka 代理都挂了。
在这种情况下,在将消息附加到其缓冲区之前,KafkaProducer 会尝试获取元数据。如果 KafkaProducer 在配置的超时时间内无法获取元数据,则会引发异常。
-内存记录不可写(kafka 0.9.0.1 库中存在错误)
https://issues.apache.org/jira/browse/KAFKA-3594
以上两种情况,KafkaProducer 都不会重试,Flink 会忽略这些消息。甚至没有记录消息。例外是,但不是失败的消息。
可能的解决方法(Kafka 设置):
- 元数据超时值非常高 (metadata.fetch.timeout.ms)
- 缓冲区过期值非常高 (request.timeout.ms)
我们仍在调查更改上述 kafka 设置可能产生的副作用。
那么,我们的理解正确吗?或者有没有办法通过修改一些 Flink 设置来避免这种数据丢失?
谢谢。
【问题讨论】:
标签: apache-kafka apache-flink producer