Конвейер данных для персонализации

Автор: Николай СапуновОбновлено: август 202619 мин чтения
Содержание статьи +

TL;DR

Конвейер данных для персонализации – это «водопровод», который превращает действия зрителей (каждый запуск, паузу, скролл) в свежие сигналы, на которых реально работают рекомендации, поиск и мерчандайзинг, и от которого тихо зависит вся остальная discovery-машина. У него четыре движущиеся части: сбор событий со всех экранов, передача через replayable-лог в stream- и batch-обработку, хранение результатов в feature store, который отдаёт моделям одни и те же числа в обучении и в проде, и граница приватности данных просмотра, которые несколько законов считают особо чувствительными. Самая дорогая инженерная ошибка – training-serving skew, когда признаки, на которых модель обучалась, расходятся с теми, что она видит вживую, и это тихо ломает рекомендации; а самая дорогая юридическая ошибка – утечка истории просмотров в сторонний тег, что в США запускает Video Privacy Protection Act с ответственностью $2,500 на человека (18 U.S.C. § 2710). Эта статья объясняет весь конвейер простым языком, чтобы вы могли решить, что строить, что покупать и где должны проходить границы.

Почему это важно

Всё остальное в discovery – ряды рекомендаций, поиск, персонализированный мерчандайзинг и эксперименты, которые их настраивают, – хорошо ровно настолько, насколько хорош питающий их конвейер данных, так что это фундамент, на котором стоит весь блок. Сделайте правильно – и выбор зрителя за последний час способен изменить домашний экран уже к вечеру; сделайте неправильно – и самая умная модель будет получать устаревшие или несогласованные сигналы, и рекомендации сгниют. Эта статья – для основателя, продакт-менеджера или CTO стриминга, которому нужно решить, сколько дата-инфраструктуры строить, должны ли признаки обновляться за секунды или за ночь и – критически – где проходит юридическая черта вокруг того, что вы знаете о зрителях. Вы не будете писать stream-код, но вам предстоит задать архитектуру и границу приватности и понять, почему самый дешёвый на вид срез угла так часто оборачивается самым дорогим багом или самым дорогим иском.

Что такое конвейер персонализации на самом деле

Начнём с самой простой формулировки. Конвейер данных для персонализации – это цепочка систем, которая берёт сырое поведение зрителей (что смотрели, когда остановились, что искали, что пролистали) и превращает его в аккуратные актуальные числа, которые читают модели рекомендаций и поиска, чтобы решить, что показать каждому человеку дальше. Это мост между беспорядочным потоком действий зрителей и чистыми входами, которые нужны алгоритму.

Полезная аналогия – кухня ресторана. Зал (ваши приложения) генерирует постоянный поток заказов и реакций; кухня (конвейер) их принимает, заготавливает ингредиенты в стандартную форму, держит самые ходовые ингредиенты наготове на линии и следит, чтобы блюдо, собранное на тесте, было таким же, как отправленное в зал. Модель рекомендаций – это шеф на раздаче, но шеф хорош ровно настолько, насколько хороша кухня за его спиной (а как именно эта модель ранжирует тайтлы – отдельное ремесло). Бо́льшая часть работы и бо́льшая часть режимов отказа живут на кухне.

Эту кухню составляют четыре задачи, и дальше каждая – отдельный раздел: собрать события со всех устройств; передать и обработать их через лог в real-time и batch вычисления; сохранить результаты в feature store, который гарантирует модели одинаковые данные в обучении и в проде; и защитить данные просмотра за границей приватности, которую закон воспринимает всерьёз. Каждая – место, где стриминговые платформы либо строят устойчивое преимущество, либо тихо ломают продукт.

Рисунок 1. Конвейер данных для персонализации от начала до конца. События со всех экранов собираются, передаются через replayable-лог, обрабатываются в двух темпах (real-time stream и плановый batch) и пишутся в feature store, который кормит модели рекомендаций, поиска и ранжирования. Пунктирная граница отмечает зону данных просмотра, где действуют правила согласия, минимизации и срока хранения.

Шаг первый: сбор событий

Всё начинается с события – небольшой записи с меткой времени, которая говорит: «этот зритель сделал это в этот момент». События – сырьё персонализации; без событий самая умная модель слепа. Они делятся на три семейства, и держать их раздельно помогает рассуждать о том, что несёт конвейер.

Первое семейство – playback-события: то, что происходит внутри плеера. Play start срабатывает в начале тайтла; pause, resume и seek (прыжок в другую точку) – по мере управления; complete – в конце. Важно, что плеер также шлёт heartbeat – небольшой пинг «ещё смотрю, на этой позиции» каждые несколько секунд, – потому что без него вы не отличите тридцатисекундный сэмпл от целого эпизода. Watch time, метрика, которая решает почти всё, восстанавливается из этих heartbeat'ов.

Второе семейство – навигационные события, иногда называемые clickstream: что зритель делает вне плеера. Показ (impression) фиксирует, что тайтл был показан в ряду (появился на экране); клик – что зритель открыл его страницу; скроллы, просмотры рядов и наведения фиксируют просмотр каталога. Показы важны не меньше кликов, потому что система рекомендаций должна учиться и на том, что она показала, а зритель проигнорировал, а не только на том, что он выбрал.

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

Эти события не приходят по одному из одного места. Один «нажал play» запускает события сразу из нескольких сервисов – playback-сервис, сервис рекомендаций, отдавший ряд, сервис мониторинга качества и сервис CDN-маршрутизации шлют свои (Netflix Technology Blog, 2018). Умножьте на каждый экран – и объём велик, что и стоит сделать конкретным.

Арифметика объёма событий

Числа задают архитектуру, поэтому пройдём по одному. Допустим, на пике у вашей платформы миллион зрителей смотрит одновременно. Плеер шлёт heartbeat каждые 10 секунд – это 6 heartbeat'ов в минуту, – и при активном просмотре каждый зритель генерирует ещё около 4 навигационных событий в минуту. Назовём это 10 событий в минуту на зрителя. Тогда:

событий/сек = 1 000 000 зрителей × 10 событий/мин ÷ 60 с
            = 10 000 000 ÷ 60
            ≈ 167 000 событий в секунду

То есть сервис с миллионом одновременных зрителей производит порядка 167 000 событий каждую секунду в этот момент – и это ровный темп, до всякого всплеска от премьеры. За целый день, даже при меньшей средней одновременности, это миллиарды событий. Для масштаба: дата-конвейер Netflix сообщает о пиках около 12,5 млн событий в секунду и порядка двух триллионов событий в день по всем сервисам (Netflix Technology Blog, 2018). Если одно сырое событие – примерно килобайт, то миллиард событий – около терабайта сырых данных, поэтому первое реальное ограничение – стоимость хранения и обработки, а не изобретательность. Это та же scale-first реальность, что управляет масштабированием и одновременностью в OTT: вы проектируете конвейер под пик, а затем удешевляете ровный режим.

Семейство событийЧто фиксируетКонкретные примерыЧто питает
PlaybackЧто происходит внутри плеераplay start, pause, resume, seek, complete, heartbeatWatch time, досмотры, «продолжить просмотр»
Навигация (clickstream)Просмотр вне плеерапоказ, клик, скролл, просмотр ряда, наведениеЧто показано против выбранного; сигналы ранжирования
ПоискНамерение, введённое зрителемзапрос, показанные результаты, выбранный результатРелевантность поиска, сигналы спроса

Таблица 1. Три семейства событий, которые собирает конвейер персонализации. Playback-события (особенно heartbeat) восстанавливают watch time; навигационные фиксируют, что показано и проигнорировано, а не только что нажато; поисковые фиксируют заявленное намерение. Модель, которая учится только на кликах и никогда на показах, не отличит хорошую рекомендацию от удачной позиции.

Шаг второй: передача и обработка событий – два темпа

После сбора события должны добраться от миллионов устройств до систем, которые считают признаки, и проектное решение, организующее всё, – это лог. Лог здесь – это append-only, replayable-запись каждого события в порядке, в котором оно произошло; думайте о нём как о плёнке, которую всегда можно перемотать. Apache Kafka – самая частая реализация. Продьюсеры (приложения и сервисы) пишут события в лог; консьюмеры (задания обработки) читают из него в своём темпе. Лог разъединяет тех, кто создаёт события, и тех, кто их использует, и, поскольку он replayable, вы можете переобработать историю, прочитав его заново с начала (Kleppmann, 2017).

Из лога обработка идёт в двух темпах, и понимание разницы – сердце проектирования конвейера.

Быстрый темп – это stream-обработка: задания, которые читают события в момент их прихода и обновляют число за секунды. «Сколько эпизодов этот зритель досмотрел за последний час?» – это streaming-признак: он должен быть актуален, чтобы быть полезным. Apache Flink – частый движок здесь. Типичное streaming-задание отфильтровывает шум, затем обогащает каждое событие – сырой play_start несёт лишь user_id, title_id и метку времени, поэтому задание джойнит его со справочными данными, чтобы прикрепить жанр, страну зрителя и возрастной рейтинг, прежде чем что-либо ниже по течению это использует (Netflix Technology Blog, 2018).

Медленный темп – это batch-обработка: задания, которые идут по расписанию – каждые несколько часов или ночью – по огромным историческим наборам. «Любимые жанры этого зрителя за последние 90 дней» – это batch-признак; он меняется медленно, так что пересчитывать его раз в день нормально и куда дешевле, чем трогать на каждом событии. Apache Spark – частый движок здесь.

У того, как вы комбинируете два темпа, есть название. Lambda-архитектура запускает оба слоя параллельно – batch-слой для точной полной истории и быстрый speed-слой для последних минут – и сливает их выходы на чтении (Marz & Warren, 2015). Она отказоустойчива, но несёт скрытый налог: вы поддерживаете одну и ту же логику дважды, в batch- и в stream-системе, а держать две кодовые базы выдающими одинаковый результат ровно так же больно, как звучит. Эта боль и привела к тому, что Джей Крепс, один из создателей Kafka, предложил в 2014-м Kappa-архитектуру: убрать отдельный batch-слой, держать единственный stream-слой над replayable-логом, а когда нужно пересчитать историю – просто проиграть лог через тот же код (Kreps, 2014). Один путь кода, один источник истины, пересчёт перемоткой плёнки.

Рисунок 2. Два темпа конвейера персонализации. Streaming-признаки должны быть свежими за секунды (эпизоды за этот час); batch-признаки меняются медленно и дёшево пересчитываются по расписанию (любимые жанры за 90 дней). Lambda держит отдельные batch- и speed-слои и сливает их; Kappa держит один stream-слой и пересчитывает проигрыванием лога – один путь кода вместо двух.

Практический совет прост. Используйте streaming-признак, когда свежесть меняет ответ: последние просмотренные тайтлы, что в тренде прямо сейчас, «продолжить просмотр». Используйте batch-признак, когда сигнал медленный, а объём огромен: долгий вкус, часы за всё время, эмбеддинги контента. Большинство реальных платформ запускают оба, и вопрос архитектуры – можете ли вы свести их к одному пути кода (Kappa) или вам реально нужны два (Lambda). Ловушка – тянуться к real-time везде: stream-инфраструктура дороже и сложнее в эксплуатации, а большинству сигналов персонализации секунды свежести не нужны.

ИзмерениеReal-time (streaming) признакиBatch-признаки
СвежестьСекундыЧасы – сутки
Типичный движокStream-процессор (напр. Flink)Batch-движок (напр. Spark)
Стоимость и эксплуатацияВыше – всегда онлайн, сложнее вестиНиже – по расписанию, проще
Применяйте для«Смотрел за час», тренды, продолжить-просмотрВкус за 90 дней, часы за всё время, медленные эмбеддинги
Отказ при неверном выбореУстаревший сигнал убивает «что в тренде»Зря потраченные деньги на ненужную свежесть

Таблица 2. Когда использовать real-time признак, а когда batch. Решающий вопрос – меняет ли свежесть рекомендацию. Выбор real-time по умолчанию везде – частая и дорогая ошибка; выбор batch там, где сигнал действительно чувствителен ко времени, заставляет «тренды сейчас» врать.

Шаг третий: feature store – и баг, который его окупает

Теперь самая важная и наименее понятая часть. Признак – это одно число, которое читает модель: «досмотрено эпизодов за неделю», «доля комедий за 30 дней». Feature store – это система, которая считает, хранит и отдаёт эти признаки, и она решает конкретную дорогую проблему, которую объясняет остаток раздела. Концепцию назвала и популяризировала платформа Uber Michelangelo в 2017-м, построившая центральный store, чтобы команды могли «создавать и управлять каноническими признаками» и переиспользовать их между моделями, а не пересобирать одну и ту же логику снова и снова (Uber Engineering, 2017).

Проблема, ради убийства которой существует feature store, – training-serving skew: разрыв между признаками, на которых модель обучалась, и признаками, которые она реально получает вживую (Zinkevich, Google «Rules of Machine Learning»). Модель обучают на исторических данных, собранных одним способом – скажем, batch-заданием по хранилищу, – а затем отдают вживую другим путём кода, который пересчитывает те же признаки под бюджет задержки. Если два пути кода считают «средний watch time за неделю» хоть чуть-чуть по-разному – иное округление, иное окно времени, иная обработка пропусков – модель видит одно в лаборатории и другое в дикой природе, и её рекомендации тихо деградируют. Урон трудно заметить именно потому, что каждый кусок по отдельности выглядит корректным; ничего не падает, числа просто расходятся.

Структурное лекарство – и есть причина существования feature store: определите каждый признак один раз и отдавайте это одно определение и в обучение, и в прод. Feature store делает это двумя синхронизированными половинами. Offline store держит длинную историю и отвечает на вопросы обучения по огромным наборам. Online store держит текущее значение каждого признака и отдаёт его за несколько миллисекунд на serving'е – Uber сообщал об отдаче online-признаков с задержкой 95-го перцентиля менее ~10 миллисекунд, достаточно быстро, чтобы поместиться внутри живого запроса рекомендации (Uber Engineering, 2017). Поскольку обе половины питаются из одного определения, модель получает согласованные числа в обоих темпах. Дополнительная защита, рекомендованная в инженерных правилах Google, – логировать ровно те признаки, что использованы на serving'е, и учить следующую модель на этих логах: если вы учитесь на том, что реально отдали, два пути не могут разъехаться (Zinkevich, Google «Rules of Machine Learning»).

Рисунок 3. Зачем нужен feature store. Два отдельных пути кода для обучения и serving'а расходятся и вызывают training-serving skew (слева, баг). Одно общее определение признака, питающее синхронизированные offline store (обучение, вся история) и online store (serving, миллисекунды), держит их согласованными (справа, решение). Point-in-time join'ы гарантируют, что строки обучения видят только данные, существовавшие до предсказываемого момента.

Есть и вторая, более тонкая ловушка, от которой защищает feature store: point-in-time корректность, она же избегание утечки данных (data leakage). Когда вы строите обучающий пример для «уйдёт ли этот зритель», вы джойните метку каждого зрителя (отменил ли он подписку) с его признаками (его поведением). Если вы случайно прикрепите значения признаков из после момента, который притворяетесь, что предсказываете, модель учится на будущем – информации, которой у неё в реальной жизни никогда не будет. Результат – модель, блестящая на тесте и проваливающаяся в проде; классический симптом – offline-точность 0,95, обваливающаяся до 0,78 вживую (Huyen, 2022). Feature store с time-travel join'ами собирает каждую обучающую строку, используя только значения признаков, существовавшие до метки времени этой строки, что нудно делать руками и легко сделать неправильно – ещё одна причина, по которой крупные стриминговые и tech-платформы построили feature store, а не выводят join заново каждый раз (Huyen, 2022).

Шаг четвёртый: граница приватности данных просмотра

Теперь конвейер знает для каждого зрителя, что именно он смотрел и когда. Это ровно та информация, которую несколько законов считают чувствительной, и граница приватности – это часть архитектуры, которая решает, какие данные куда, кому и как надолго утекают. Относиться к ней как к запоздалой мысли – так победа персонализации становится коллективным иском.

Специфичный для стриминга закон, который нужно знать, – американский Video Privacy Protection Act (VPPA), кодифицированный в 18 U.S.C. § 2710. Он принят в 1988-м после того, как газета опубликовала историю видеопроката кандидата в Верховный суд, и его базовое правило прямолинейно: «video tape service provider» не может сознательно раскрывать личностно идентифицируемую информацию – определённую в законе как информацию, идентифицирующую человека «как запросившего или получившего конкретные видеоматериалы или услуги» – без информированного письменного согласия зрителя (18 U.S.C. § 2710(a)(3), (b)). Зубы – в ущербе: суд может присудить ликвидированный ущерб $2,500 на каждого пострадавшего (§ 2710(c)(2)), что в коллективном иске на миллионы зрителей огромно. Хотя закон написан для видеопрокатов, язык применяется к современному стримингу, и недавняя волна исков целит ровно в один паттерн – отправку событий просмотра сторонним рекламным и аналитическим тегам, – почему он и стоит в центре главной ловушки этой статьи.

Ещё две обязанности из того же закона прямо формируют конвейер. Согласие по VPPA должно быть конкретным и отдельным – не зашитым в общий клик по условиям использования – и может быть дано заранее на срок до двух лет с понятным способом его отозвать (§ 2710(b)(2)(B)). И закон требует от провайдеров уничтожать личностно идентифицируемую информацию «как только это практически осуществимо, но не позднее одного года» после того, как она больше не нужна для своей цели (§ 2710(e)) – жёсткий лимит хранения, который политика хранения вашего конвейера должна реально обеспечивать, а не просто обещать.

Помимо VPPA, применимы два широких режима, потому что данные просмотра, привязанные к аккаунту или устройству, – это персональные данные. Европейский GDPR (Регламент (ЕС) 2016/679) требует законного основания для их обработки (статья 6) и задаёт принципы, которые должна воплощать ваша архитектура: ограничение цели (используйте только для того, о чём сказали зрителю), минимизация данных (собирайте только нужное) и ограничение срока хранения (храните не дольше необходимого) – статья 5(1)(b), (c), (e). Калифорнийский CCPA, с поправками CPRA (Cal. Civ. Code § 1798.100 и далее), даёт зрителям права знать, удалять и отказываться от «продажи» или «передачи» их данных и считает точные данные о потребителе чувствительными. Работа с приватностью, разобранная сквозным образом в статье приватность и данные просмотра: VPPA, GDPR, CCPA, – юридический спутник этой инженерной статьи; здесь же суть в том, что конвейер – это место, где эти правила соблюдаются или нарушаются.

На практике граница строится из четырёх контролей. Псевдонимизация заменяет реальную личность непрозрачным ключом внутри конвейера, чтобы аналитические системы работали с viewer_8f3c1, а не с именем и почтой. Минимизация данных означает не собирать поля, которым нет применения. Лимиты хранения удаляют сырые события по часам, удовлетворяющим годовому правилу VPPA и ограничению срока хранения GDPR. А состояние согласия путешествует вместе с данными, чтобы у зрителя, который не соглашался на рекламный таргетинг, история просмотров никогда не утекла рекламному партнёру. Граница – это не одна стена; это набор правил, прикреплённых к данным по мере их движения.

Рисунок 4. Граница приватности данных просмотра. Внутри границы события просмотра псевдонимизированы, минимизированы, ограничены по сроку и гейтятся состоянием согласия. Красный путь – передача личности зрителя плюс конкретного просмотренного тайтла стороннему рекламному или аналитическому тегу – это ровно тот паттерн, что запускает ответственность по Video Privacy Protection Act $2,500 на человека (18 U.S.C. § 2710).
РежимЧто регулируетКлючевая обязанность для конвейераКонтроль, который он навязывает
VPPA (18 U.S.C. § 2710)Раскрытие кто что смотрелНет раскрытия PII просмотра без конкретного письменного согласия; уничтожить ≤ 1 года после конца нуждыБлок сторонних тегов; гейт согласия; часы хранения
GDPR (ЕС 2016/679)Любые перс. данные зрителей ЕСЗаконное основание; ограничение цели, минимизация, срок (ст. 5–6)Минимизировать поля; удалять по расписанию; задокументировать цель
CCPA / CPRA (Калиф.)Перс. данные жителей КалифорнииПрава знать, удалять, отказаться от продажи/передачиИсполнять удаление/opt-out; отделять «переданные» данные

Таблица 3. Три режима приватности, наиболее релевантных конвейеру персонализации, и контроль, который каждый навязывает в архитектуру. VPPA – специфичный для стриминга; его ликвидированный ущерб на человека делает утечку в сторонний тег самой высокорисковой одиночной ошибкой во всём конвейере.

Частая ошибка: удобный сторонний тег

Ловушка, превращающая конвейер персонализации в иск, – заодно и самая безобидная на вид. Команда хочет лучшую аналитику или атрибуцию рекламы, поэтому ставит сторонний тег – рекламный пиксель или аналитический SDK – в приложение и даёт ему наблюдать события просмотра. Это одна строка интеграции, и она «просто работает». Проблема в том, что это ровно тот акт, который запрещает VPPA: провайдер раскрывает внешней стороне информацию, идентифицирующую конкретного человека как смотревшего конкретный контент (18 U.S.C. § 2710(b)). Волна коллективных исков по VPPA за последние несколько лет целила ровно в этот паттерн, а порог закона $2,500 на человека означает, что даже скромная база пользователей подразумевает очень крупную экспозицию.

Инженерное лекарство – граница приватности выше: события просмотра никогда не покидают границу привязанными к личности, если согласие явно не разрешает, сторонние теги держат подальше от потока просмотра, а то, что вы всё же передаёте, псевдонимизировано и минимизировано. Вторая, более тихая версия «удобного среза угла» – на стороне моделирования: считать признак одним способом для обучения и другим для serving'а, потому что написать второй запрос было быстрее, чем переиспользовать одно определение. Обе ошибки растут из одного инстинкта: относиться к конвейеру как к клею, а не как к продукту. На стриминговой платформе конвейер данных и есть продукт персонализации.

Где здесь Фора Софт

Конвейер персонализации – это задача масштаба и корректности раньше, чем задача машинного обучения: модель – малая часть, а устойчивое преимущество – кухня за её спиной: сбор событий, не теряющий данные на пике, feature store, отдающий одинаковые числа в обучении и проде, и держащая граница приватности. За 250+ реализованными проектами для 400+ клиентов с 2005 года в видеостриминге, OTT/Internet TV, e-learning и телемедицине мы строим полный конвейер: сбор событий на web, mobile и TV; replayable-лог, питающий и stream-, и batch-обработку; feature store, устраняющий training-serving skew и обеспечивающий point-in-time корректность; и граница приватности, спроектированная вокруг VPPA, GDPR и CCPA с первого события, а не прикрученная после запуска. Наш подход scale-first и vendor-neutral: мы стартуем от вашего пикового объёма событий и того, какие сигналы реально должны быть свежими до секунд, решаем, где хватит хостингового feature store и stream-сервиса, а где кастомный конвейер окупает свою цену, и заводим выход прямо в системы рекомендаций, поиска и экспериментов, чтобы то, что зритель сделал за последний час, могло изменить то, что он увидит вечером, – безопасно.

Ключевые выводы

  • Конвейер персонализации – собрать, передать, хранить, защитить – фундамент, на котором стоит каждая модель рекомендаций и поиска.
  • Собирайте три семейства событий: playback (heartbeat восстанавливает watch time), навигацию и поиск.
  • Real-time признаки – только когда свежесть меняет ответ; batch дешевле для медленных сигналов.
  • Feature store существует, чтобы убить training-serving skew: одно определение признака служит и обучению, и проду.
  • Point-in-time join'ы останавливают утечку – строки обучения не должны видеть данные из будущего относительно момента.
  • Данные просмотра юридически чувствительны: никогда не сливайте «кто-что-смотрел» сторонним тегам (VPPA, $2,500 на человека).

Что почитать дальше

Строите такую систему?

Подберём параметры кодирования под ваш контент и посчитаем стоимость доставки до старта разработки.