The Jerusalem PostIsraeli soccer team to wear black armbands against Ireland in honor of October 7 victimsRTP DesportoPortugal goleia Índia no arranque do Mundial feminino de hóquei em patinsPunchNorthern senators seek military action against Adamawa insurgentsESPNSources: LaPorta, Lions break off contract talks after negotiations stallESPN DeportesColts y Commanders elevan el nivel de sus ofensivas y se cierra el juegoInquirerWalang pasok: In-person classes suspended Oct. 5 for Nat’l Teachers’ DayZDF heuteEntdecken Sie das ZDF-NachrichtenstudioSouth China Morning PostHong Kong builders take three WorldSkills honours in best CIC haulColliderTom Cruise’s Divisive New Movie Officially Ends His Box Office Winning StreakVarietyJohn Steinbeck’s ‘East of Eden’ Returns to Top of Book Charts Following Netflix Adaptation’s ReleaseHipertextualCómo prestar tu móvil sin miedo: el truco oculto de Android que nadie usaSRF NewsGülsha über Leben und Musik – Ein Tantra-Retreat brachte Gülsha zum Weinen
The Daily Newsstand · Free, Always
Sunday, October 4, 2026

chunk() пропустил половину рассылки

Translate

Раз в минуту у нас просыпаются четыре консольные команды. Каждая проходит по таблице чатов — это сотни тысяч строк и около трёхсот тысяч живых пользователей в месяц — и решает, кому бот напишет первым: у кого обновились суточные лимиты, кто начал разговор и замолчал на третьей минуте, у кого закончилась подписка, кому пора показать промо.

Год назад я полез разбираться, почему уведомление об обновлении лимитов приходит с опозданием на несколько минут, хотя команда отрабатывает за секунды и в логе чисто. Оказалось, что за один проход команда обслуживает половину тех, кого сама же выбрала. Не «примерно половину» — ровно половину. Виноват chunk(), который отработал в точности так, как написано в документации.

Проект под NDA, поэтому дальше без имён: продуктовая специфика убрана, сущности переименованы, бизнес-цифры не называются. Механика и грабли настоящие.

Что бот пишет сам

У телеграм-бота есть неочевидное свойство: он умеет начинать разговор. Пользователь не открывает приложение — приложение приходит к нему само. Это одновременно главный канал возврата и главный способ получить блокировку.

Инициативных сообщений у нас четыре вида:

  • обновились суточные лимиты — «заходи, у тебя снова есть чем пользоваться»;

  • человек начал разговор и не написал ни одного сообщения — через три минуты ему приходит второе сообщение от собеседника;

  • закончилась подписка — предложение продлить;

  • давно не заходил — промо.

Каждое — отдельная команда, и все четыре стоят на everyMinute().

Schedule::command('chat:premium-finished')->everyMinute()->withoutOverlapping();
Schedule::command('chat:second-message')->everyMinute()->withoutOverlapping()->runInBackground();
Schedule::command('chat:refill')->everyMinute()->withoutOverlapping()->runInBackground();
Schedule::command('chat:promo')->everyMinute()->withoutOverlapping()->runInBackground();

Почему раз в минуту, а не раз в сутки

Первая версия обновления лимитов была ночным кроном: в полночь пройтись по всем и всем всё выдать. Так делают почти все, и на маленькой базе это работает.

Ломается с двух сторон сразу. Со стороны нагрузки — в полночь ты получаешь один большой UPDATE по всей таблице и пачку из десятков тысяч сообщений, которые надо отправить в течение нескольких минут, потому что «лимиты обновились» через час уже никому не интересно. Со стороны продукта — полночь у всех разная, а у нас пять локалей и пользователи по всем часовым поясам.

Поэтому лимиты стали скользящим окном от активности: у каждого чата своя метка, через сутки после которой ему полагается обновление. Крон перестал быть раздатчиком и стал догонялкой — он ищет тех, кто не заходил дольше суток, обновляет им лимиты и зовёт обратно. Тех, кто зашёл сам, обслуживает вебхук по дороге, без всякого крона.

Побочный эффект приятный: вместо пика на всю базу получается ровный поток. Побочный эффект неприятный: запросы, которые раньше выполнялись раз в сутки и никого не волновали, теперь выполняются 1440 раз в сутки — и вот тут выясняется, как они на самом деле написаны.

Как chunk() читает таблицу

Команда обновления лимитов выглядела так — сокращённо, без локалей и логов:

Chat::query()
    ->join('billings', 'chats.id', '=', 'billings.chat_id')
    ->whereNotNull('billings.tokens_reset_at')
    ->where('billings.tokens_reset_at', '<=', $limit)
    ->where('chats.is_blocked', false)
    ->whereNull('chats.banned_at')
    ->orderBy('billings.tokens_reset_at', 'asc')
    ->select('chats.*')
    ->chunk(120, function ($chats) use ($service) {
        foreach ($chats as $chat) {
            $data = $service->performRefill($chat);   // внутри: tokens_reset_at = null
            $this->sendRefillNotification($chat, $data);
        }
    });

performRefill внутри транзакции выдаёт лимиты и ставит tokens_reset_at = null — это метка «цикл закрыт, новый начнётся при следующей активности». То есть обработанная строка перестаёт подходить под условие выборки. Ровно в этом и проблема.

chunk() — это не курсор. Это цикл, который каждую итерацию выполняет отдельный запрос с LIMIT и OFFSET:

-- первая порция
select chats.* from chats join billings ... where billings.tokens_reset_at <= ? ... limit 120 offset 0
-- вторая порция
select chats.* from chats join billings ... where billings.tokens_reset_at <= ? ... limit 120 offset 120

Между двумя запросами мы обнулили метку у первых ста двадцати строк. Ко второму запросу они уже не проходят по условию — выборка сдвинулась на 120 строк влево. А offset 120 отсчитывает от нового начала. Первая порция забрала строки 1–120, вторая забирает 241–360, строки 121–240 не увидит никто.

Дальше по индукции: обслуживается половина, аккуратными чередующимися полосами по 120 строк.

Сколько это стоило на самом деле

Первая реакция была «ну и ладно, следующей минутой догонит». Так и есть: через минуту крон стартует с нуля, сортировка по времени метки ставит пропущенных в начало, и они получают своё.

Но посчитаем. Тысяча чатов в очереди на обновление — это не тысяча за проход, а 500, потом 250, потом 125. Чтобы разгрести тысячу, нужно около десяти минут вместо одной. Пока очередь короткая, это невидимо. В день, когда очередь стала длинной, это стало выглядеть как «уведомления приходят с задержкой» — за эту ниточку я и дёрнул.

Хуже другое: в логе чисто. Команда не падает, каждая обработанная строка честно пишет «отправлено», метрика растёт. Дырка между «сколько подходило под условие» и «сколько обработали» не была видна нигде, потому что первое число никто не считал.

Как chunk() теряет строки: выборка сдвигается между запросами

Как chunk() теряет строки: выборка сдвигается между запросами

chunkById и почему он не вставляется в одну строчку

Лечится заменой на chunkById(), который вместо OFFSET тащит курсор по первичному ключу:

select ... where billings.tokens_reset_at <= ? and chats.id > ? order by chats.id asc limit 120

Строки, выпавшие из выборки, больше не сдвигают окно: следующий запрос начинается с конкретного id, а не с «отступи 120 от начала». Приём известен как keyset pagination, и он же лечит вторую, менее заметную болезнь OFFSET — растущую стоимость на больших смещениях.

Две вещи, о которых узнаёшь уже на замене.

Первая: chunkById() переопределяет сортировку. Мой orderBy по времени метки был не украшением, а смыслом — «сначала те, кто ждёт дольше всех». После перехода порядок стал по chats.id, то есть по дате регистрации: свежие пользователи начали ждать за спинами тех, кто зарегистрировался два года назад. Справедливость очереди пришлось возвращать иначе — ограничивать верхнюю границу метки и брать пачку целиком, а не полагаться на ORDER BY внутри обхода.

Вторая: с джойном колонку курсора нужно называть полностью, иначе база не поймёт, чей id имеется в виду:

->chunkById(120, function ($chats) { /* ... */ }, 'chats.id', 'id');

Третий аргумент — колонка в запросе, четвёртый — имя ключа в полученной модели. Если перепутать, получишь либо Column 'id' in where clause is ambiguous, либо, что веселее, бесконечный цикл: курсор будет читать не ту колонку и никогда не дойдёт до конца.

Надёжнее: сначала id, потом работа

chunkById чинит сдвиг окна, но не чинит главного: мы по-прежнему читаем и пишем одну и ту же выборку в одном проходе. Там, где проход длинный, я теперь делаю иначе — сначала собираю идентификаторы, потом работаю:

$ids = Chat::query()
    ->join('billings', 'chats.id', '=', 'billings.chat_id')
    ->where('billings.tokens_reset_at', '<=', $limit)
    // ... остальные условия
    ->orderBy('billings.tokens_reset_at')
    ->limit(self::BATCH)
    ->pluck('chats.id');

Выборка зафиксирована, порядок сохранён, а лимит на пачку появился сам собой — и это внезапно самое ценное. Крон, который раз в минуту берёт не больше N чатов, деградирует предсказуемо: при завале растёт задержка, а не время прохода. Крон без такого лимита при завале начинает наезжать сам на себя.

Цена очевидная: между pluck и обработкой строка может измениться, и обрабатывать её уже не надо. Проверять актуальность всё равно придётся, но теперь это дешёвая проверка перед отправкой, а не невидимый пропуск в середине обхода.

С другой стороны — дубликаты

Команда второго сообщения устроена иначе: она не отправляет сама, а ставит задачи в очередь.

->chunk(120, function ($chats) {
    foreach ($chats as $chat) {
        SendSecondMessageJob::dispatch($chat->id)->onQueue('second_message');
    }
});

Здесь тот же сдвиг выборки — джоба выставляет second_message_sent_at, и строка выпадает из условия. Но добавляется беда, зеркальная первой: метка выставляется не сразу, а когда до задачи дойдут руки воркера. Если очередь отстаёт на полторы минуты, следующий запуск крона увидит те же чаты и поставит те же задачи ещё раз.

В самой джобе защита есть:

if ($this->chat->second_message_sent_at !== null) {
    return;
}

Это фильтр, а не защита. Он спасает от последовательного выполнения дублей и не спасает от параллельного: два воркера берут две копии задачи, оба читают null, оба отправляют. Пользователь получает два одинаковых сообщения подряд — выглядит ровно так, как выглядит баг.

Честных вариантов три, и все дешёвые:

  1. ShouldBeUnique на джобе с ключом по идентификатору чата — Laravel возьмёт лок в Redis на время выполнения;

  2. пометить строку в момент диспатча: update ... where second_message_sent_at is null и смотреть на количество затронутых строк;

  3. уникальный индекс на таблице отправленных и ловля нарушения — так у нас сделаны резервы баланса, про них была отдельная статья.

Мне больше нравится второй: метка ставится тем же запросом, который её проверяет, и рассинхрон между «выбрали» и «пометили» исчезает как класс — без локов и без дополнительной инфраструктуры. У нас пока стоит только проверка внутри джобы, то есть первый вариант из трёх не реализован, а третий применён в другом месте системы.

withoutOverlapping() и сутки тишины

Все четыре команды стоят с withoutOverlapping(), и это правильно: проход по сотням тысяч строк может не уложиться в минуту, и накладываться ему нельзя.

Чего я не знал: у лока есть время жизни, и по умолчанию оно — 24 часа. Лок живёт в кэше, снимается в конце выполнения, и если процесс умер не своей смертью — OOM-killer, kill -9, перезагрузка сервера в неудачный момент, — снимать его некому. Команда молча не выполняется. Ровно сутки.

У нас это случилось один раз, с командой про закончившуюся подписку. Обнаружилось не по мониторингу, а по тому, что за день не пришло ни одного продления.

Schedule::command('chat:premium-finished')->everyMinute()->withoutOverlapping(5);

Пять минут — это «в несколько раз больше, чем самый долгий легальный проход». Дальше остаётся вторая половина проблемы: после суток простоя команда просыпается и видит не десяток строк, а тысячи. А там:

$billings = Billing::where('premium_until', '<', now())->with('telegramChat')->get();

get() без всяких порций. В нормальном режиме это десяток строк в минуту, и жило оно так годами. После суток тишины это память, которой нет. Чинится одной строкой — и, как обычно, чинилось уже после.

Отправка внутри крона

Две команды из четырёх до сих пор отправляют сообщения синхронно, прямо в процессе крона. Выглядит невинно:

foreach ($chats as $chat) {
    $this->sendPromo($chat, $promos->random());
}

Внутри — HTTPS-запрос к Telegram API. Даже при быстрых ответах это сотня-другая миллисекунд на чат, последовательно, в одном процессе. Сто двадцать чатов — полминуты. Тысяча — четыре минуты при минутном расписании, и withoutOverlapping() эти запуски просто съест.

Правильный ответ — крон только выбирает и ставит задачи, отправляют воркеры. У нас так сделана одна команда из четырёх, и это долг, а не архитектурное решение. Отдельная причина сделать именно так: в очереди отправка уже обложена ретраями, приоритетами и лимитами. В кроне соблюдать лимиты Telegram попросту нечем — там нет ни общего счётчика, ни места, где его держать.

Мёртвые души

Любая выборка для рассылки начинается не с того, кого мы хотим позвать, а с того, кого звать нельзя:

->where('chats.is_blocked', false)
->whereNull('chats.banned_at')
->where('chats.is_started', true)

is_blocked ставим не мы, а Telegram: когда пользователь блокирует бота, API на любую отправку отвечает 403. Это единственный способ узнать о блокировке — попробовать написать.

case 403:
    $this->chat->is_blocked = true;
    $this->chat->save();
    break;

Дальше начинается арифметика, которая на старте проекта кажется неважной. Доля заблокировавших только растёт: тот, кто заблокировал бота, не разблокирует его никогда. Через год это заметная часть таблицы, и если выборка её не отсекает, крон каждую минуту честно обходит кладбище — с запросами, с попытками отправки, с местом в очереди, которое могло достаться живому.

Поэтому индекс под эти три колонки появился раньше, чем индекс под саму метку времени.

Индекс, который пришлось перевернуть

Индексов под эти выборки два:

$table->index(['bot_id', 'is_blocked', 'banned_at']);   // chats
$table->index(['chat_id', 'tokens_reset_at']);          // billings

Второй изначально был написан наоборот — ['tokens_reset_at', 'chat_id']. Логика казалась очевидной: фильтруем по времени, значит время первым.

Она была бы верной, будь billings ведущей таблицей в плане. Но план начинается с chats — там отсекается всё лишнее, заблокированные и забаненные, — и в billings мы приходим уже за конкретным chat_id. При таком плане индекс с временем впереди не используется вовсе: ведущая колонка в соединении не участвует.

Миграция короче, чем объяснение:

Schema::table('billings', function (Blueprint $table) {
    $table->dropIndex(['tokens_reset_at', 'chat_id']);
    $table->index(['chat_id', 'tokens_reset_at']);
});

Порядок колонок в составном индексе определяется планом запроса, а не важностью колонок в голове автора. EXPLAIN до и после занимает две минуты, и обе эти минуты я в тот раз пожалел.

Чего не хватало в мониторинге

Всё описанное выше — один класс ошибок: расхождение между «сколько строк подходило под условие» и «сколько мы обработали». Ни одна из них не ловится алёртом на ошибки, потому что ошибок нет.

Что из этого следует для мониторинга:

  • каждая команда в конце должна писать три числа: сколько нашла, сколько обработала, сколько пропустила осознанно;

  • если «нашла» больше суммы двух других — алёрт, независимо от причины;

  • длина очереди на обслуживание — сколько строк подходит под условие прямо сейчас — отдельная метрика с графиком.

Третий пункт, подозреваю, полезнее первых двух: он показывает проблему до того, как её заметит пользователь. У фоновых рассылок вообще нет естественного индикатора здоровья — они по определению работают без человека, который пожалуется.

Сейчас у нас из этих трёх есть только счётчик отправленных, то есть самое бесполезное. Дырку между «нашли» и «обработали» я в своё время нашёл руками, из любопытства, и это худший из возможных способов.

Что я меняю после этой статьи

Пока писал, перечитал все четыре команды подряд, чего давно не делал. Список того, что поеду чинить:

  1. chunk() на движущейся выборке остался ещё в двух командах — там же, где и был.

  2. Дубликаты задач: метку надо ставить в момент диспатча.

  3. withoutOverlapping() без явного времени жизни — везде.

  4. get() без порций в команде про подписку.

  5. Синхронная отправка в двух командах из четырёх.

Ни один из пяти пунктов не проявляется на тестовой базе в тысячу строк. Все пять проявляются на живой.

Чеклист

Если у вас есть фоновая команда, которая ходит по большой таблице и что-то в ней меняет:

  1. chunk() нельзя использовать, если обработка выводит строки из выборки. Только chunkById() или явная пачка идентификаторов.

  2. chunkById() переопределяет сортировку — если порядок был содержательным, его надо возвращать другим способом.

  3. С джойном указывайте колонку курсора полностью, вместе с именем таблицы.

  4. Лимит на размер прохода нужен всегда: при завале должна расти задержка, а не длительность прохода.

  5. Помечайте строку тем же запросом, который её выбирает, а не позже и не в другом процессе.

  6. Джоба, идемпотентная только проверкой «если уже сделано — выходим», не идемпотентна.

  7. withoutOverlapping() — всегда с явным временем жизни, кратным самому долгому легальному проходу.

  8. Команда, которая нормально живёт на десяти строках в минуту, обязана пережить сутки простоя. Проверяется руками: остановить, подождать, запустить.

  9. Никаких синхронных сетевых вызовов в цикле крона. Крон выбирает, очередь отправляет.

  10. Выборка для рассылки начинается с исключений, а не с условий. Заблокированные, забаненные, не стартовавшие — первыми.

  11. Порядок колонок в составном индексе проверяется EXPLAIN, а не рассуждением.

  12. Логируйте «нашли / обработали», а не «отправлено». Расхождение этих двух чисел — единственный способ увидеть тихий пропуск.

Отдельно любопытно про пятый пункт. Мы ставим метку в момент диспатча и миримся с тем, что при падении воркера пользователь не получит сообщение вовсе. Обратный вариант — метка после успешной отправки — гарантирует доставку ценой дублей. Третьего мы не придумали, а хочется: как вы разруливаете эту развилку у себя?

View the original on Хабр →

KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.