【问题标题】:Kafka Connect for Elastic search fails to publish few messages out of 10million messages to ElasticKafka Connect for Elastic search 未能将 1000 万条消息中的几条消息发布到 Elastic
【发布时间】:2021-09-25 17:12:00
【问题描述】:

我们设置了一个 Kafka 连接,用于将事件从 Kafka 发送到 Elastic。我们遇到了数据一致性问题,希望在这里获得一些指导。

我们使用的是开源的kafka connect confluentinc/cp-kafka-connect-base:6.1.0 它对我们的日常使用非常有用,但通过负载测试我们发现了问题。 我们应用的负载在 Kafka 主题上产生了 1000 万个独特事件。当我们检查弹性搜索时,我们发现丢失了 96 个文档,它们从未进入弹性搜索。

我们检查了 kafka 连接日志和运行 kafka 连接的 kubernetes pod,发现如下:

  1. pod 上的 CPU 非常高,进而导致 pod 在负载测试的整个生命周期中多次被杀死。不确定这是否可能是此丢失文档的问题。我们现在通过添加更多 CPU 和自动缩放解决了这个问题。
  2. 在日志中,我们看到弹性服务本身也返回了网关超时错误 (504),这似乎又导致批量请求失败,并出现如下错误:
    ERROR WorkerSinkTask{id=kafka-connect-gds-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. Error: Bulk request failed. (org.apache.kafka.connect.runtime.WorkerSinkTask)
  1. 没有其他错误,DLQ主题中也没有消息。

  2. 我们手动检查并比较了使其具有弹性的事件与没有发现任何差异的事件,因此我们认为此网关超时错误一定是导致 96 个文件无法显示的问题松紧带。但我们可能错了。

在网关错误消失并重新启动 Pod 后,此问题也应该得到解决。计数仍然没有。所以看起来 SinkTask 未能写入弹性,但提交了消费者组的偏移量?这就是为什么重启后 Kafka 连接没有接收到那些从未到达 Elastic 的事件。 我希望这两件事可以在交易中完成。我试图在这里阅读实现https://github.com/apache/kafka/blob/trunk/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java,但无法真正找出正确的流程。

所以我希望能在这里得到一些关于我们可能做错了什么的见解。

这是我们的 Kafka 连接设置:

{
    "name": "${KC_ES_CONNECTOR_NAME}",
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "connection.url": "${KC_ES_DOMAIN_ENDPOINT}",
    "type.name": "kafkaconnect",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "schema.ignore": "true",
    "topics": "${KC_ES_TOPICS}",
    "behavior.on.null.values": "DELETE",
    "behavior.on.malformed.documents": "IGNORE",
    "batch.size": "500",
    "read.timeout.ms": "60000",
    "linger.ms": "100",
    "write.method": "upsert",
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "${KC_ES_NAMESPACE}.kc.gds.dlq",
    "errors.deadletterqueue.context.headers.enable": "true",
    "errors.log.enable": "true",
    "errors.log.include.messages": "true",
    "errors.retry.timeout": "-1",
    "errors.retry.delay.max.ms": "60000",
    "predicates": "isNullRecord",
    "predicates.isNullRecord.type": "org.apache.kafka.connect.transforms.predicates.RecordIsTombstone",
    "transforms": "dropNullRecords,transformserecords",
    "transforms.dropNullRecords.type": "org.apache.kafka.connect.transforms.Filter",
    "transforms.dropNullRecords.predicate": "isNullRecord",
    "transforms.transformserecords.type": "com.mycompany.kafka.connect.smt.TransformSERecord\$Value",
    "transforms.transformserecords.transform": "true"
}

我们也有一个自定义 SMT,它会进行非常基本的转换,例如添加一些新字段并忽略某些事件并按照约定将记录路由到正确的弹性索引。所有 10M 文件都是相同的,因此转换对所有文件的工作方式都相同,我们可以通过 SMT 的日志记录确认这一点。

发生此问题时,Kafka 连接正在使用 1 个工作人员和 1 个 pod(进程)。绝对没有为测试类型提供准备,但我们已经解决了这个问题。我们计划重复测试,但希望就我们在这里遗漏的任何明显内容提供指导。

问候, 维卡斯

【问题讨论】:

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


    【解决方案1】:

    Kafka connect 不会接收那些从未进入 Elastic 的事件。我希望这 2 件事可以在事务中完成。

    在 kafka 内部贡献者社区中有关于在接收器/源连接器上支持更精确一次行为的实时和持续讨论,很难实现这种行为,因为两个系统之间没有共享事务,大多数连接器至少支持一次或进行 upsert 以创建幂等操作。

    特别是关于 Confluent 的弹性搜索接收器连接器

    一次交货 连接器依赖 Elasticsearch 的幂等写入语义来确保准确地一次交付给 Elasticsearch。通过在 Elasticsearch 文档中设置 ID,连接器可以确保只发送一次。如果 Kafka 消息中包含密钥,则它们会自动转换为 Elasticsearch 文档 ID。当 key 不包含,或者被明确忽略时,connector 会使用 topic+partition+offset 作为 key,确保 Kafka 中的每条消息在 Elasticsearch 中都有一个对应的文档。

    死信队列 此连接器支持死信队列 (DLQ) 功能。有关访问和使用 DLQ 的信息,请参阅 Confluent 平台死信队列。

    【讨论】:

    • 感谢 Ran 您的意见。因此,基本上您在差距上所描述的内容是不可能保证所有 Kafka 事件都可以写入弹性文件。如果由于某种原因弹性不可用,事件将像我们在测试中看到的那样被丢弃。我们的活动中确实有唯一的 ID,如果 ES 一直在运行,则不会出现交付问题。
    猜你喜欢
    • 2019-09-22
    • 2018-06-04
    • 2019-09-01
    • 1970-01-01
    • 2020-11-26
    • 2019-01-22
    • 2019-11-18
    • 2020-04-27
    • 1970-01-01
    相关资源
    最近更新 更多