Vocabulary for Stream Processing: Windowing, Watermarks, and Backpressure

Освоєння англійської мови для обробки потоків: вікно, водяний знак, зворотній тиск, exactly- once, пізні дані і оператори стану. Точні терміни для інженерів з обробки даних і платформ.

Обробка потоків — Flink, Kafka Streams, Spark Structured Streaming — має щільний, спеціалізований словник. Такі слова, як windowing, watermark і backpressure не зустрічаються в повсякденній англійській, тому немає інтуїції, на яку можна покластися. У цьому довіднику визначено основні терміни, показано колокації дієслів, які використовують справжні інженери, і позначено пастки у вимові. Освоивши їх, ви зможете обговорювати потокові канали з точністю.


Потоки проти пакетів: основна відмінність

Задача ** пакета ** обробляє фіксовану, обмежену групу даних і завершує роботу. ** Потік ** обробляє ** необмежений **, постійний потік подій, який ніколи не закінчується.

  • ** Обмежені проти необмежених** дані
  • ** У спокійному стані проти у руху ** — дані, що знаходяться у сховищі проти потоку даних через конвеєр
  • ** Обробка в реальному часі / майже в реальному часі **
  • ** Прохідність ** (кількість подій за секунду) і ** затримка ** (час на подію)

«Ми перенесли агрегування з нічного ** пакетного ** завдання на ** потокове ** конвеєр, тому панелі управління оновлюють в ** майже реальному часі ** замість одного разу на день.»


Час події проти часу обробки

Це розрізнення є основою обробки потоків і найскладнішим для новачків:

  • ** Час події ** — час, коли подія * фактично сталася * (часовий штамп у даних).
  • ** Час обробки ** — час, коли ваша система * отримала* його.
  • Время поглощения — когда оно вошло в трубопровод.

Ці параметри відрізняються через затримки мережі, повторні спроби і автономні пристрої. Фраза, яку ти постійно будеш чути:

«Ми агрегуємо за часом події, а не часом обробки, тому що мобільні клієнти буферизують події офлайн і надсилають їх пізно»


Windowing

Оскільки потік є нескінченним, ви не можете « чекати на всі дані ». Замість цього ви групуєте події у ** вікна ** — кінцеві шматки для обчислення.

  • ** Вікно з перемикачами ** — фіксоване, не перекривається (наприклад, кожні 5 хвилин).
  • ** Слайд- вікно ** — перекривається (наприклад, останні 5 хвилин, оновлюється кожну хвилину).
  • ** Вікно сеансу ** — групується за діяльністю, закривається після перерви у бездіяльності.
  • ** Розмір вікна ** і ** інтервал слайдів **

Колокацій:

  • події ** впадають ** у вікно
  • вікно ** запускає ** / ** закриває ** / ** видає результат **
  • ми ** розширимо поток на ** 5 хвилин
  • вікно ** вмикає ** обчислення

“Ми використовуємо ** вікно падіння ** одну хвилину. Коли вікно загоряється, воно видає рахунок і очищає свій стан

** Вимова: ** * TUM- bling, * скольжение * (SLY-ding). Прямо, але скажи їх з впевненістю.


Водяні знаки і пізні дані

** Водяний знак ** це оцінювання системи « ми, ймовірно, бачили всі події до часу T. » Це дозволяє рушію вирішувати, коли вікно є достатньо повним для виведення.

  • ** Водяний знак ** — рухомий поріг повноти часу події.
  • ** Пізні дані / події з пізнім прибуттям ** — події, які з’ являться * після * того, як їх вікно пройде водяний знак.
  • ** Дозволена затримка ** — період відстрочки для запізнених подій.
  • ** Відкинуті / відкинуті ** події — події, які відбулися пізніше за період відстрочки.

Водяний знак збільшується з часом події. Події, що прибувають після водяного знака, є затримками; в межах дозволеної затримки ми все ще рахуємо їх, за межами цього ми викидаємо їх»

Це класичне питання для інтерв’ ю з працівниками потокових послуг — вам слід мати змогу пояснити водяний знак у двох реченнях.


Backpressure

** Backpressure ** це те, що відбувається, коли оператор нижче не може слідкувати за оператором вище — система сигналізує « сповільнити », щоб не закінчилася пам’ ять.

  • повільний споживач ** використовує протидію ** / ** відштовхує назад **
  • трубопровід ** застосовує протитиск ** щоб зменшити джерело
  • без неї система відстає і затримка зростає
  • ми буферизуємо події, а потім починаємо дренувати продюсера

“Ринок не міг триматися, тому він ** використовував протитиск ** весь шлях назад до джерела Кафки, який ** дросел ** поглинання. Lag stayed bounded — саме те, що ми хочемо»

Вимова: одне слово, зворотний тиск (ЗВОРТ-прес-ер). Не говори “зворотний тиск” як два наголошені слова в технічному контексті.


Обробка станів і станів

Оператори потоків часто потребують state — запам’ятовування речей між подіями.

  • ** Stateless ** — кожна подія обробляється незалежно (наприклад, фільтр).
  • ** Stateful ** — залежить від минулих подій (наприклад, від поточного підрахунку).
  • ** Keyed state ** — стан розділений за допомогою ключа.
  • ** Сервер стану ** — місцезнаходження стану (пам’ ять, RocksDB).
  • ** Checkpoint ** — періодичний знімок стану для відновлення.
  • ** Savepoint ** — знімок, який створюється вручну для оновлення.

“Агрегація є статевою: вона зберігає загальну суму в ключовому стані, підтримувану RocksDB, і контрольними точками кожні 30 секунд, щоб вона могла відновлюватися після аварії.”


Семантика доставки

Ті ж три гарантії, що і для повідомлень, але тут вони критично важливі:

  • ** Майже раз ** — може втратити події.
  • ** Принаймні- раз ** — може бути дублікатом.
  • ** Точно- один раз ** — кожна подія впливає на стан один раз (золотий стандарт, важко досягти).

«Ми покладаємося на checkpoints і idempotent sinks, щоб досягти exactly-once семантики від початку до кінця»

Порівняння: idempotent sink, transactional writes, replay з джерела при відновленні.


Когда что-то пойдет не так

  • Lag — наскільки далеко відстає конвеєр від реального часу. “Consumer lag is climbing.”
  • ** Схилення ** — нерівномірний розподіл даних, що перевантажує один розділ (** схилення даних**, ** скорочення клавіш**).
  • ** Затримки ** — повільні завдання, які затримують решту.
  • ** Переповнення диска ** — стан занадто великий для пам’ яті.

Гаряча клавіша спричинила зсув даних — один оператор зробив 80% роботи, в той час як інші працювали на паузі. Ми переключилися на ребаланс»


До / після: вільно говорить

** До: ** « Іноді події надходять пізно, і підрахунок неправильний, і коли наступний крок повільний, пам’ ять заповнюється всім »

** Після: ** « **Пізнє надходження подій ** за водний знак було відкинуто, що спричинило викривлення підрахунку. І коли стік сповільнився, відсутність зворотного тиску спричинила розлив стану розлив і вибух пам’яті. “


Швидкий словник

TermOne-line meaning
UnboundedInfinite, never-ending stream
Event timeWhen the event actually happened
WindowA finite chunk of the stream to compute over
WatermarkEstimate of event-time completeness
Late dataEvents arriving after the watermark
BackpressureDownstream signalling “slow down”
CheckpointSnapshot of state for recovery
Exactly-onceEach event affects state once
Skew / hot keyUneven load on one partition

Ключевые вещи

  • Потік є ** необмеженим ** — ви обчислюєте через ** вікна **, а не весь набір даних.
  • Розрізняти час події від час обробки; агрегувати за часом події для правильності.
  • ** Водяні знаки ** визначають, коли вікно завершено; події, що відбуваються після них, є ** запізненими ** і можуть бути ** відкинуті **.
  • ** Backpressure ** утримує швидкого виробника від перевантаження повільного споживача — дізнайтеся колокацію “**exert backpressure **.”
  • Конвейери з відомим станом покладаються на контрольні точки для відновлення і неможливі ями для точно-одиничного доставлення.

Навигація нюансів: спільне виразування та відгуки в обговореннях обробки потоків

Будьмо чесними – коли обговорюються складні технічні концепції, такі як обробка потоків, точна мова має багато значення. Це не тільки про розуміння “чого”, але і “як” ти це передаєш. Для не-рідних носіїв англійської мови це може бути особливо викликом, особливо коли справа доходить до спеціалізованого жаргону. Просте нерозуміння терміну або фрази може призвести до плутанини під час перегляду коду, обговорення з зацікавленими сторонами або навіть зневадження проблем.

Однією з спільних областей, де виникають нюанси, є якість даних і обробка пізніх даних. Під час нещодавнього перегляду коду нової програми Kafka Streams, яка обробляє дані потоку клацань веб- сайту, рецензент залишив коментар: « Розгляньте можливість реалізації стратегії * пізнього вікна * — ми бачимо, що події надходять значно пізніше визначеного періоду часу вікна, що впливає на наші показники точності ». Це не просто говорить про те, що « дані надходять пізно ». Це означає, що потрібно прийняти певне рішення щодо архітектури програми, пов’ язане з тим, як програма обробляє ці запізнені події. Більш пряма фраза, наприклад, «Дані запізнилися», не передала б технічні наслідки або не запропонувала б потенційне рішення. Аналогічно, при описі потреби в коректуванні водяного знака в описі PR, краще сказати «Ми повинні збільшити водяний знак, щоб враховувати збільшення кількості подій», ніж просто сказати «Водяний знак потребує коригування»

Інша часто зустрічається фраза — «точно-один разова семантика», яка часто викликає плутанину. Це не про отримання * точно * ті ж дані назад; це про те, щоб переконатися, що кожен запис обробляється * не більше одного разу *, незалежно від невдач або повторних спроб. Повідомлення Slack під час сеансу усунення несправностей виглядало приблизно так: «Гаразд, давайте зосередимося на ідемпотентних операторах і транзакційних гарантіях, щоб досягти обробки точно-один раз. Нам потрібно переконатися, що стан не пошкоджений дублікатними оновленнями, якщо споживач не впорався.” Зауважте наголос на гарантіях - це про встановлення довіри до поведінки системи. Також корисно розуміти, що досягнення «все одноразового» часто включає в себе складні конфігурації навколо транзакційних можливостей Kafka і ідемпотентних ключів.

І, нарешті, не бійтеся просити про пояснення. Якщо ви чуєте термін, який ви не до кінця розумієте, завжди краще ввічливо запитати. Правильним правилом є конструктивне формулювання вашого питання: « Чи могли б ви розібратися, що ви маєте на увазі під « вікном » у цьому контексті? » Це демонструє вашу зацікавленість і бажання вчитися, які є високоцінними рисами у будь- якому інженерному середовищі.

Ось простий приклад використання вікнового оператора time у Kafka Streams для ілюстрації налаштування водяного знака (за допомогою Java):

import org.apache.kafka.streams.state.*;
import java.util.*;

public class TimeWindowExample {

    public static void main(String[] args) {
        // This is a simplified example - real-world configurations are more complex
        int windowDurationSeconds = 5; // Define the window duration (e.g., 5 seconds)
        long watermarkOffset = 0L;  // Initial watermark offset

        System.out.println("Time Window Duration: " + windowDurationSeconds + " seconds");
        System.out.println("Initial Watermark Offset: " + watermarkOffset);
    }
}

Цей приклад демонструє основну концепцію - * тривалість вікна * (5 секунд в цьому випадку) визначає період, за який події агрегуються, а watermarkOffset відстежує позицію пізніх даних в цьому вікні. Налаштування водяного знаку є ключовим способом обробки даних, які не збігаються з послідовністю або затримуються.

Поширені запитання

Про що ця стаття "Vocabulary for Stream Processing: Windowing, Watermarks, and Backpressure"?

Освоєння англійської мови для обробки потоків: вікно, водяний знак, зворотній тиск, exactly- once, пізні дані і оператори стану. Точні терміни для інженерів з обробки даних і платформ.

Чи безкоштовна ця стаття?

Так. Усі статті на CoderSlingo, включно з цією, доступні безкоштовно без реєстрації.

Скільки часу займає читання "Vocabulary for Stream Processing: Windowing, Watermarks, and Backpressure"?

Приблизно 11 min.