AbsaOSS / AbsaOSS/hyperdrive

Add deduplication logic for kafka sink in case of spark task retries

未关闭
#240 0 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看
enhancement
主要语言
Scala
星标
47
派生
14
PR 合并指标
30 天内没有已合并 PR

描述

**Problem description**
Spark does not provide an exactly-once behaviour for the Kafka sink, but only at-least-once, and will probably never do so (https://github.com/apache/spark/pull/25618). Under certain assumptions (no concurrent producers, only 1 destination topic, not too big micro-batches, messages don't change between retries), idempotency can still be achieved. See #177.

The DeduplicateKafkaSinkTransformer only addresses retries on an application level. However, retries (and therefore duplicates) may happen on lower levels as well, namely:

- Retry of a `DataWritingSparkTask` (the Spark task that will invoke the KafkaProducer)
- Internal retry of the `KafkaProducer` (when encountering a `RetriableException`, e.g. server disconnected)
Duplicates due to an internal retry of the `KafkaProducer` can be prevented by setting `acks=all` and `enable.idempotence` on the kafka writer. However, this does not take into account retries of the `DataWritingSparkTask` which invokes the `KafkaProducer`. For example, if there is an intermittent `TopicAuthorizationException`, the `KafkaProducer` will fail and not retry, but the `DataWritingSparkTask` will retry nevertheless. In such a case, duplications are still possible. Another example is executor failure due to exceeding memory limits. If an executor exceeds memory limits during the `DataWritingSparkTask`, it will be terminated and another executor will retry the task, which may again lead to duplicates on the destination topic.
One solution is to switch off spark task retries, but obviously, this causes other problems.

**Solution**
It might be possible to implement a `org.apache.kafka.clients.producer.ProducerInterceptor` which would deduplicate retries messages in a similar way as the `DeduplicateKafkaSinkTransformer`. This would capture duplicates in a scenario where the retry happens on a different executor than the original executor (memory exceeded scenario). In addition, the interceptor should keep a lookup set in memory to capture duplicates that occurred due to retries on the same executor
(TopicAuthorizationException scenario).
A reattempt could be recognized through the attempt nr of the Spark task.
If messages cannot be dropped through the ProducerInterceptor, they could at least be redirected to a trash-topic

贡献指南

这个仓库没有索引到贡献指南

评估

这个 Issue 还没有评估数据。

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。