Англійська для розробників Apache Beam
Словник для розробників, які створюють конвеєри Apache Beam — PCollections, windowing, watermarks і transforms — для команд, які обговорюють обробку пакетів і потоків англійською мовою.
Вся концепція Beam полягає в тому, що пакетна обробка є лише особливим випадком потокової передачі, і більшість її словників існує для опису того, як вона обробляє дані, які постійно надходять з часом — включаючи дані, які надходять пізно. Зрозуміти цей словник - це те, що робить потокові огляди трубопроводів зрозумілими, а не залякаючими.
Основні напрямки
** Pipeline ** — повний графік перетворень даних, які виконує Beam, переносимий між запусковими програмами (Dataflow, Flink, Spark) без переписування самої логіки.
- “Ми написали цей конвеєр один раз і він працює на Dataflow в prod і на прямому локальному runner — ця портативність є головною метою використання Beam.” *
** PCollection ** — основна абстракція даних Beam: незмінна, можливо необмежена збірка елементів, на яких працюють перетворення, відрізняється від звичайного списку у пам’ яті, оскільки може представляти нескінченний потік.
“Ви не можете просто викликати .length на PCollection — вона може бути необмежена, тому Beam не має поняття про кінцевий розмір, поки ви не обмежите її.”
** Перетворення (PTransform) ** — операція, яка приймає одну або декілька PCollections як вхідні дані і створює одну або декілька PCollections як вихідні дані, що є складовою частиною конвеєра Beam.
- “Загорнути цю логіку у власну PTransform, щоб її можна було використовувати у трьох конвеєрах, які потребують одного кроку дедуплікації.” *
Час і вікна
** Розгортання ** — розділення необмеженого PCollection на скінченні шматки за часом, таким чином, агрегації (наприклад, підрахунки або суми) мають обмежений, значущий обсяг для роботи.
- “Без вікон, « кількість подій за хвилину » навіть не має сенсу — Beam має знати, де закінчується одна хвилина і починається наступна.” *
** Водяний знак ** — оцінювання Beam щодо того, наскільки повними є дані до певного часу події, використовується для визначення того, коли вікно можна вважати « завершеним », навіть якщо дані все ще можуть надходити пізно.
- “Водяний знак пройшов за межі вікна, отже Beam запустив агрегування — але ми все ще дозволяємо пізнім даним запускати оновлення результату.” *
** Пізні дані ** — елементи, які надходять після того, як водяний знак вже пройшов час події, що вимагає явного правила (відкинути або викликати виправлення), а не мовного оброблення.
- “Це число змінилося після факту через пізні дані — ми дозволили вікно виправлення, тому агрегований результат було перераховано, як тільки з’ явилась пізня подія.” *
** Тригер ** — правило, яке визначає, коли програма Beam надсилатиме результати для вікна, наприклад, один раз, коли пройде водяний знак, або кілька разів, коли буде надсилатися більше даних.
- “Ми додали ранній тригер, щоб панелі отримували приблизну кількість кожної хвилини, замість очікування, поки вікно повністю закривається.” *
Поширені помилки
- Припустимо, що вікно агрегату є остаточним у момент його виведення, без урахування тригерів пізніх даних, які можуть переглянути його пізніше.
- Розгляд « необмежених » як синонімів « пошкоджених » або « з вадами » замість навмисної, правильної поведінки потокової PCollection.
- Забувши, що переносимість конвеєра між забігачами не означає ідентичних характеристик продуктивності - часто все ще потрібний спеціальний тунінговий прохід.
Практичні вправи
- Поясніть у двох реченнях, чому вимога щодо віконності є необхідною перед тим, як ви зможете зрозуміло « підрахувати події за хвилину » у потоці.
- Напишіть короткий коментар щодо перегляду дизайну, у якому поясніть, яким чином можна обмежити водяний знак і дозволити виправлення, які надходять пізніше.
- Виконати чернетку повідомлення, у якому буде пояснено, що « необмежений » розмір PCollection є очікуваною поведінкою, а не звітом про помилку.
Зв’язані ресурси
- Англійська для розробників Apache Spark
- Англійська для розробників Kafka Streaming
- Англійська для розробників Apache Flink
Навігація по зворотному зв’язку та налаштування Windows
Початкове захоплення отримання роботи Beam pipeline часто уступає місце більш нюансованим завданням його вдосконалення - особливо при роботі з потоками даних. Велика частина цього включає отримання та відповідь на зворотній зв’ язок, не тільки з автоматичних перевірок, але також з обговорень вашої команди. Часто мова, що використовується навколо вікон і водяних знаків, може бути досить технічною, що призводить до нерозуміння, якщо не ретельно сформулювати. Як людина, для якої англійська не є рідною, я вважаю, що важливо зосередитись на ясному спілкуванні про намір за цими конфігураціями, а не занурюватися в точні математичні визначення. Наприклад, рецензент може сказати: « Цей розмір вікна здається агресивним; чи ви впевнені, що ви захоплюєте всі відповідні події? » Хороша відповідь буде не просто « Я використовував вікно з тривалістю 10 секунд », а щось на зразок « Ми обирали вікно з тривалістю 10 секунд, щоб переконатися, що ми захопили пік активності під час цієї акції. Ми також додали водяний знак, щоб запобігти необмеженому зростанню, тому навіть якщо швидкість перевищує цей поріг, нові дані не будуть включені у вікно. ”
Інший поширений сценарій виникає при обговоренні водяних знаків - особливо їх вплив на повність. Повідомлення Slack може звучати так: « Чи хтось хвилюється про втрату даних через водяний знак? » Знову ж таки, точна технічна відповідь не є негайно корисною. Замість цього, пояснення логіки є ключовим. «Водяний знак, який ми встановили на 10 МБ на вікно, розроблений для запобігання надмірних витрат на розрахунки і забезпечення того, що ми не зберігаємо застарілих даних, які більше не відображають поточні умови. Ми ретельно стежимо за показниками повноти і, якщо буде потрібно, вносити зміни на основі спостережуваних тенденцій. Важливо проактивно вказати * чому * було обрано певне значення водяного знака — оптимізація витрат, вимоги до відповідності або просто практичні обмеження зберігання.
Крім того, при описі складних перетворень в конвеєрі Beam, корисно описати їх з точки зору їх впливу на бізнес. Замість того, щоб сказати « Ми використовували перетворення CombineWindowed з агрегацією sumByKey », що може звучати надто технічно, ви можете сказати: « Це перетворення агрегує дані про продажі за категоріями продуктів протягом 5-хвилинних вікон, що дозволяє нам швидко визначити періоди пікового попиту і оптимізувати рівень запасів ». Цей підхід перетинає прогалину між технічною реалізацією і кінцевою метою конвеєра.
Ось приклад простого налаштування трубопроводу Beam за допомогою Python:
import apache_beam as beam
def run():
with beam.Pipeline() as pipeline:
lines = pipeline | 'Create' >> beam.Create(['a', 'b', 'c', 'd'])
squared = lines | 'Square' >> beam.Map(lambda x: x * x)
results = squared | 'Format' >> beam.Map(str)
print(results)
if __name__ == '__main__':
run()