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

Год назад я полез разбираться, почему уведомление об обновлении лимитов приходит с опозданием на несколько минут, хотя команда отрабатывает за секунды и в логе чисто. Оказалось, что за один проход команда обслуживает половину тех, кого сама же выбрала. Не «примерно половину» — ровно половину. Виноват 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. Логируйте «нашли / обработали», а не «отправлено». Расхождение этих двух чисел — единственный способ увидеть тихий пропуск.

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