【问题标题】:Kafka Consumer with DLQ and ElasticSearch带有 DLQ 和 ElasticSearch 的 Kafka 消费者
【发布时间】:2019-03-24 08:12:59
【问题描述】:

我有以下集群:

Kafka -> 一些日志收集器 -> Elasticsearch

我的问题是关于选择最高效的日志收集器(或其他一些允许管理 Kafka 和 ElasticSearch 之间的数据流的软件)。

我正在尝试从 LogstashFluentdConfluent 的 Kafka Elasticsearch 连接器中进行选择。 我面临的主要问题是在写入 Elasticsearch 端点时出现问题后无法在 Kafka 中回滚偏移量。

例如,logstash 文档说“400404 错误会发送到死信队列 (DLQ)(如果启用)。如果未启用 DLQ,将发出一条日志消息,并且该事件将被丢弃”(https://www.elastic.co/guide/en/logstash/6.x/plugins-outputs-elasticsearch.html#_retry_policy)。如果我有这样的错误,logstash 将继续从 Kafka 读取数据。错误会一次又一次地发生。虽然,我所有的数据都将存储到 DLQ 中,但当第一个错误发生时,Kafka 的偏移量将远离该位置。我必须手动定义正确的偏移量。

所以,我的问题是: Kafka 和 ElasticSearch 是否有任何连接器,它允许在收到来自 ElasticSearch (400/404) 的第一个错误后停止移动偏移量?

提前致谢。

【问题讨论】:

    标签: elasticsearch apache-kafka logstash apache-kafka-connect fluentd


    【解决方案1】:

    我认为问题不在于效率,而在于可靠性

    我面临的主要问题是在写入 Elasticsearch 端点出现问题后,无法在 Kafka 中回滚偏移量。

    我对 Connect 或 Logstash 的 DLQ 功能没有太多经验,但重置消费者组偏移量并非不可能。但是,如果消费者应用程序正确处理偏移提交,则不需要这样做。

    如果 Connect 向 ES 抛出连接错误,它将重试,而不是提交偏移量。

    如果错误不可恢复,则 Connect 将停止消费,并且再次不提交偏移量。

    因此,从消息批处理中获取丢失数据的唯一方法是,如果该批处理以 DLQ 结尾,使用任何框架。

    如果禁用 DLQ,丢失数据的唯一方法是从 Kafka 过期

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-25
      • 2014-07-30
      • 1970-01-01
      • 2018-04-28
      • 2018-04-02
      • 1970-01-01
      • 2020-09-19
      • 1970-01-01
      相关资源
      最近更新 更多