【发布时间】: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,发现如下:
- pod 上的 CPU 非常高,进而导致 pod 在负载测试的整个生命周期中多次被杀死。不确定这是否可能是此丢失文档的问题。我们现在通过添加更多 CPU 和自动缩放解决了这个问题。
- 在日志中,我们看到弹性服务本身也返回了网关超时错误 (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)
-
没有其他错误,DLQ主题中也没有消息。
-
我们手动检查并比较了使其具有弹性的事件与没有发现任何差异的事件,因此我们认为此网关超时错误一定是导致 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