Add deduplication logic for kafka sink in case of spark task retries
Nadie ha tomado este issue todavía.
Evaluación
- Dificultad
- 5/5
- Tiempo estimado
- Más de una semana
- Aptitud para principiantes
- 25/100
- Tipo de issue
- Nueva funcionalidad
- Claridad
- Bastante claro
- Estado de actividad
- Estancado
- Stack tecnológico
- kafka, scala
- Área
- stream-processing
Línea de trabajo
Comienza revisando el DeduplicateKafkaSinkTransformer existente y el número de intento de la tarea de Spark disponible durante los retries. Evalúa un Kafka ProducerInterceptor que gestione los retries entre executors y mantenga una búsqueda en memoria para los retries en el mismo executor; se considera terminado cuando los duplicados se descartan o se redirigen a un trash topic bajo las suposiciones indicadas.
Escrito por el modelo de indexación a partir del texto del issue.
Descripción
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
- Lenguaje dominante
- Scala
- Estrellas
- 47
- Forks
- 14
- Métricas de merge de PR
- Sin PR fusionados en 30 d
Preparar el entorno
Este proyecto no incluye contenedor de desarrollo, Dockerfile ni guía de contribución, así que la configuración corre por tu cuenta: empieza por su README y consulta nuestra guía para la primera contribución para los pasos generales.
Primeros pasos
- Lee el issue completo y luego la guía de contribución del proyecto.
- Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
- Haz un fork del repositorio y trabaja en una rama.
- Abre un pull request que haga referencia al número del issue.
Más de AbsaOSS/hyperdrive
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 45/100
AbsaOSS/hyperdrive#269 ·
-
enhancement
Dificultad 3/5 1-2 días Aptitud para principiantes 42/100
AbsaOSS/hyperdrive#239 ·
-
bug
Dificultad 3/5 1-2 días Aptitud para principiantes 35/100
AbsaOSS/hyperdrive#230 ·
-
Atum integrationQuizá libre de nuevo @kevinwallimann la tomó hace 2005 días y no hay ningún pull request abierto. Abierto
Dificultad 5/5 Más de una semana Aptitud para principiantes 20/100
AbsaOSS/hyperdrive#211 · 1 comentario · 1 reacción · 1 asignado ·
-
Improve robustness of component testsQuizá libre de nuevo @kevinwallimann la tomó hace 2289 días y no hay ningún pull request abierto. Abiertointernal-task
Dificultad 4/5 3-5 días Aptitud para principiantes 28/100
AbsaOSS/hyperdrive#148 · 1 asignado ·
Todos los issues de AbsaOSS/hyperdrive
Issues similares
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 78/100
Los mantenedores suelen responder en 1 día
-
[Rust][Flaky Test] multiple_deadlines_fire_in_order asserts a wall-clock gap instead of firing orderAbiertoCI/CD ⚒️ Flaky-tests 🐦
Dificultad 2/5 1-3 horas Aptitud para principiantes 88/100
valkey-io/valkey-glide#7255 ·
Los mantenedores suelen responder en 3 días
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
lichess-org/lila#21905 ·
Los mantenedores suelen responder en 1 día
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
Los mantenedores suelen responder en 1 día
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 72/100
Los mantenedores suelen responder en 3 días