Add deduplication logic for kafka sink in case of spark task retries
まだ誰も着手していません。
評価
- 難易度
- 5/5
- 見積もり時間
- 1週間以上
- 初心者へのやさしさ
- 25/100
- issue の種類
- 機能追加
- 明瞭さ
- おおむね明確
- 活発さ
- 停滞
- 技術スタック
- kafka, scala
調査の方向性
まず、既存の DeduplicateKafkaSinkTransformer と、リトライ中に利用できる Spark タスクの試行番号を確認します。executor 間のリトライを処理し、同じ executor 内でのリトライ用にインメモリ検索を維持する Kafka ProducerInterceptor を評価します。指定された前提の下で、重複が破棄されるか trash topic にリダイレクトされれば完了です。
索引モデルが issue の本文から書いたものです。
説明
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 aRetriableException, e.g. server disconnected)
Duplicates due to an internal retry of theKafkaProducercan be prevented by settingacks=allandenable.idempotenceon the kafka writer. However, this does not take into account retries of theDataWritingSparkTaskwhich invokes theKafkaProducer. For example, if there is an intermittentTopicAuthorizationException, theKafkaProducerwill fail and not retry, but theDataWritingSparkTaskwill 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 theDataWritingSparkTask, 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
- 主要言語
- Scala
- スター
- 47
- フォーク
- 14
- PR マージ指標
- 30日以内にマージされた PR はありません
環境構築
このプロジェクトには開発コンテナ、Dockerfile、コントリビューションガイドがありません。まず README を読み、一般的な手順ははじめてのコントリビューションガイドを参照してください。
はじめの一歩
- issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
- 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
- リポジトリをフォークし、ブランチを切って変更します。
- issue 番号を参照したプルリクエストを送ります。
AbsaOSS/hyperdrive のほかの issue
-
難易度 2/5 1〜3時間 初心者へのやさしさ 45/100
AbsaOSS/hyperdrive#269 ·
-
enhancement
難易度 3/5 1〜2日 初心者へのやさしさ 42/100
AbsaOSS/hyperdrive#239 ·
-
bug
難易度 3/5 1〜2日 初心者へのやさしさ 35/100
AbsaOSS/hyperdrive#230 ·
-
Atum integration再び着手できるかも @kevinwallimann が 2009 日前に担当しましたが、オープン中のプルリクエストはありません。 オープン
難易度 5/5 1週間以上 初心者へのやさしさ 20/100
AbsaOSS/hyperdrive#211 · コメント 1 件 · リアクション 1 件 · 担当者 1 名 ·
-
Improve robustness of component tests再び着手できるかも @kevinwallimann が 2292 日前に担当しましたが、オープン中のプルリクエストはありません。 オープンinternal-task
難易度 4/5 3〜5日 初心者へのやさしさ 28/100
AbsaOSS/hyperdrive#148 · 担当者 1 名 ·
AbsaOSS/hyperdrive の issue をすべて見る
似ている issue
-
難易度 2/5 1〜3時間 初心者へのやさしさ 65/100
com-lihaoyi/mill#7670 ·
メンテナーはふだん 1 日以内に返信
-
area:aggregation bug priority:medium
難易度 2/5 1〜3時間 初心者へのやさしさ 82/100
apache/datafusion-comet#6661 ·
メンテナーはふだん 1 日以内に返信
-
難易度 2/5 1〜3時間 初心者へのやさしさ 78/100
snowflakedb/spark-snowflake#673 ·
-
bug
難易度 2/5 1〜3時間 初心者へのやさしさ 86/100
salesforce/evalon#16 ·
-
難易度 2/5 1〜3時間 初心者へのやさしさ 72/100
chipsalliance/chisel#5504 ·
メンテナーはふだん 1 日以内に返信