Вступить в клуб →
сложнаяpythonПотоковая поставка: Kafka, CDC и окна

Разбор пачки сообщений с повторами и отправкой в DLQ

Напишите process_batch(messages, handler, max_retries).

  • messages — список словарей {"offset": int, "value": str} одной партиции по возрастанию смещения.
  • handler(value) возвращает результат или бросает исключение. TransientError — временная ошибка: повторить вызов, всего не больше 1 + max_retries попыток. Любое другое исключение — постоянное: сразу в DLQ.
  • Если временная ошибка не прошла за все попытки — остановиться: сообщение не обрабатывается, последующие тоже.
  • Верните словарь: "results" — список результатов успешных сообщений по порядку; "dlq" — список (offset, имя класса исключения); "commit" — смещение, которое нужно закоммитить: следующее после последнего обработанного или отправленного в DLQ сообщения. Если ничего не обработано — None.

Класс TransientError объявлен в заготовке — не удаляйте его.

заготовка решения

class TransientError(Exception):
    pass


def process_batch(messages, handler, max_retries):
    # ваш код
    return None

🔒 Проверка ответа — для участников клуба

  • Автотесты прямо в браузере
  • Подсказка, если застряли
  • Разбор с объяснением, почему так
  • Прогресс по всем задачам и виртуальные собеседования
Зарегистрироваться →

Регистрация занимает минуту

Другие задачи раздела