Всем привет. Делимся переводом заключительной части статьи, подготовленной специально для студентов курса . С первой частью можно ознакомиться .
Apache Beam dhe DataFlow për konvej të kohës reale

Настройка Google Cloud
Примечание: Для запуска конвейера и публикации данных пользовательского лога я использовал Google Cloud Shell, поскольку у меня возникли проблемы с запуском конвейера на Python 3. Google Cloud Shell использует Python 2, который лучше согласуется с Apache Beam.
Чтобы запустить конвейер, нам нужно немного покопаться в настройках. Тем из вас, кто раньше не пользовался GCP, необходимо выполнить следующие 6 шагов, приведенных на этой .
После этого нам нужно будет загрузить наши скрипты в облачное хранилище Google и скопировать их в нашу Google Cloud Shel. Загрузка в облачное хранилище достаточно тривиальна (описание можно найти ). Чтобы скопировать наши файлы, мы можем открыть Google Cloud Shel из панели инструментов, щелкнув первый значок слева на рисунке 2 ниже.

Figura 2
Команды, которые нам нужны для копирования файлов и установки необходимых библиотек, перечислены ниже.
# Copy file from cloud storage
gsutil cp gs://<YOUR-BUCKET>/ * .
sudo pip install apache-beam[gcp] oauth2client==3.0.0
sudo pip install -U pip
sudo pip install Faker==1.0.2
# Environment variables
BUCKET=<YOUR-BUCKET>
PROJECT=<YOUR-PROJECT>Создание нашей базы данных и таблицы
После того, как мы выполнили все шаги, связанные с настройкой, следующее, что нам нужно сделать, это создать набор данных и таблицу в BigQuery. Есть несколько способов сделать это, но самый простой — использовать консоль Google Cloud, сначала создав набор данных. Вы можете выполнить действия, указанные по следующей , чтобы создать таблицу со схемой. Наша таблица будет иметь 7 столбцов, соответствующих компонентам каждого пользовательского лога. Для удобства мы определим все столбцы как строки (тип string), за исключением переменной timelocal, и назовем их в соответствии с переменными, которые мы сгенерировали ранее. Схема нашей таблицы должна выглядеть как на рисунке 3.

Рисунок 3. Схема таблицы
Публикация данных пользовательского лога
Pub/Sub является критически важным компонентом нашего конвейера, поскольку позволяет нескольким независимым приложениям взаимодействовать друг с другом. В частности, он работает как посредник, позволяющий нам отправлять и получать сообщения между приложениями. Первое, что нам нужно сделать, это создать тему (topic). Достаточно просто перейти в Pub/Sub в консоли и нажать CREATE TOPIC.
Приведенный ниже код вызывает наш скрипт для генерации данных лога, определенных выше, а затем подключается и отправляет журналы в Pub/Sub. Единственное, что нам нужно сделать, — это создать объект PublisherClient, указать путь к теме с помощью метода topic_path и вызвать функцию publish me topic_path и данными. Обратите внимание, что мы импортируем generate_log_line из нашего скрипта stream_logs, поэтому убедитесь, что эти файлы находятся в одной папке, иначе вы получите ошибку импорта. Затем мы можем запустить это через нашу google-консоль, используя:
python publish.pyfrom stream_logs import generate_log_line
import logging
from google.cloud import pubsub_v1
import random
import time
PROJECT_ID="user-logs-237110"
TOPIC = "userlogs"
publisher = pubsub_v1.PublisherClient()
topic_path = publisher.topic_path(PROJECT_ID, TOPIC)
def publish(publisher, topic, message):
data = message.encode('utf-8')
return publisher.publish(topic_path, data = data)
def callback(message_future):
# When timeout is unspecified, the exception method waits indefinitely.
if message_future.exception(timeout=30):
print('Publishing message on {} threw an Exception {}.'.format(
topic_name, message_future.exception()))
else:
print(message_future.result())
if __name__ == '__main__':
while True:
line = generate_log_line()
print(line)
message_future = publish(publisher, topic_path, line)
message_future.add_done_callback(callback)
sleep_time = random.choice(range(1, 3, 1))
time.sleep(sleep_time)Как только файл запустится, мы сможем наблюдать вывод данных лога на консоль, как показано на рисунке ниже. Этот скрипт будет работать до тех пор, пока мы не используем CTRL+C, чтобы завершить его.

Рисунок 4. Вывод publish_logs.py
Написание кода нашего конвейера
Теперь, когда мы все подготовили, мы можем приступить к самой интересной части — написанию кода нашего конвейера, используя Beam и Python. Чтобы создать Beam-конвейер, нам нужно создать объект конвейера (p). После того как мы создали объект конвейера, мы можем применить несколько функций одну за другой, используя оператор pipe (|). В общем, рабочий процесс выглядит как на рисунке ниже.
[Final Output PCollection] = ([Initial Input PCollection] | [First Transform]
| [Second Transform]
| [Third Transform]) В нашем коде мы создадим две пользовательские функции. Функцию regex_clean, которая сканирует данные и извлекает соответствующую строку на основе списка PATTERNS, используя функцию re.search. Функция возвращает разделенную запятыми строку. Если вы не являетесь экспертом по регулярным выражениям, я рекомендую ознакомится с этим и попрактиковаться в блокноте, чтобы проверить код. После этого мы определяем пользовательскую ParDo-функцию под названием Split, e cila është një variacion i transformimit Beam për përpunimin paralel. Në Python, kjo bëhet në një mënyrë të veçantë — na duhet të krijojmë një klasë që trashëgon klasën DoFn Beam. Funksioni Split merr një varg të shparshuar nga funksioni i mëparshëm dhe kthen një listë dictionaries me çelësa që korrespondojnë me emrat e kolonave në tabelën tonë të BigQuery. Ka disa gjëra që duhet të theksohen për këtë funksion: pata nevojë të importoja datetime brenda funksionit, në mënyrë që të funksiononte. Merrja një mesazh gabimi kur provoja të importoja në fillim të files, gjë që ishte e çuditshme. Kjo listë më pas kalon në funksionin WriteToBigQuery, i cili thjesht shton të dhënat tona në tabelë. Kodi për Batch DataFlow Job dhe Streaming DataFlow Job është më poshtë. E vetmja ndryshim midis kodit të paketuar dhe atij streaming është se në përpunimin për paketim ne lexojmë CSV nga src_path, duke përdorur funksionin ReadFromText nga Beam.
Batch DataFlow Job (përpunimi i paketave)
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from google.cloud import bigquery
import re
import logging
import sys
PROJECT='user-logs-237110'
schema = 'remote_addr:STRING, timelocal:STRING, request_type:STRING, status:STRING, body_bytes_sent:STRING, http_referer:STRING, http_user_agent:STRING'
src_path = "user_log_fileC.txt"
def regex_clean(data):
PATTERNS = [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(?> beam.io.textio.ReadFromText(src_path)
| "Pastrim adresash" >> beam.Map(regex_clean)
| 'Parso CSV' >> beam.ParDo(Split())
| 'Shkruaj në BigQuery' >> beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJECT), schema=schema,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
)
p.run()
if __name__ == '__main__':
logger = logging.getLogger().setLevel(logging.INFO)
main()Streaming DataFlow Job (përpunimi i rrjedhës)
from apache_beam.options.pipeline_options import PipelineOptions
from google.cloud import pubsub_v1
from google.cloud import bigquery
import apache_beam as beam
import logging
import argparse
import sys
import re
PROJECT="user-logs-237110"
schema = 'remote_addr:STRING, timelocal:STRING, request_type:STRING, status:STRING, body_bytes_sent:STRING, http_referer:STRING, http_user_agent:STRING'
TOPIC = "projects/user-logs-237110/topics/userlogs"
def regex_clean(data):
PATTERNS = [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(?> beam.io.ReadFromPubSub(topic=TOPIC).with_output_types(bytes)
| "Dekodosh" >> beam.Map(lambda x: x.decode('utf-8'))
| "Pastrim të Dhënave" >> beam.Map(regex_clean)
| 'Parso CSV' >> beam.ParDo(Split())
| 'Shkruaj në BigQuery' >> beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJECT), schema=schema,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
)
result = p.run()
result.wait_until_finish()
if __name__ == '__main__':
logger = logging.getLogger().setLevel(logging.INFO)
main()
Ekzekutimi i konvejerit
Ne mund të eciw një pipeline në disa mënyra të ndryshme. Nëse do të donim, mund të thjesht e lëshonim lokal nga terminali, duke u regjistruar nga larg në GCP.
python -m main_pipeline_stream.py
--input_topic "projects/user-logs-237110/topics/userlogs"
--streamingMegjithatë, ne do ta lëshojmë atë duke përdorur DataFlow. Mund ta bëjmë këtë me komandën e mëposhtme, duke vendosur parametrat e detyrueshëm.
project— ID e projektit tuaj GCP.runner— mjeti që ekzekuton pipeline-në, i cili analizon programin tuaj dhe ndënjti pipeline-in tuaj. Për ekzekutimin në cloud keni nevojë të specifikoni DataflowRunner.staging_location— rruga në depo të ruajtjes Cloud Dataflow për indeksimin e paketave të kodit të nevojshme për procesorët që ekzekutojnë punën.temp_location— rruga në depo të ruajtjes Cloud Dataflow për ruajtjen e skedave përkohësore të punëve, të krijuara gjatë punës së pipeline-it.streaming
python main_pipeline_stream.py
--runner DataFlow
--project $PROJECT
--temp_location $BUCKET/tmp
--staging_location $BUCKET/staging
--streaming
Ndërsa ky komandë po ekzekutohet, ne mund të kalojmë në skedën DataFlow në konsolën google dhe të shikojmë tubin tonë. Duke klikuar në tub, duhet të shohim diçka të ngjashme me figurën 4. Për qëllime debugimi, mund të jetë shumë e dobishme të kalosh te log-et, dhe pastaj në Stackdriver për të parë log-et e detajuara. Kjo më ndihmoi të zgjidh problemet me tubin në disa raste.

Figura 4: Tubi Beam
Q 접근하기 데이터 뱅크 BigQuery
Pra, ne duhet të kemi një tub të aktivizuar me të dhëna që po hyjnë në tabelën tonë. Për ta verifikuar këtë, mund të kalojmë në BigQuery dhe të shikojmë të dhënat. Pas përdorimit të komandës më poshtë, duhet të shihni disa rreshta të parë të grupit të të dhënave. Tani që kemi të dhëna të ruajtura në BigQuery, ne mund të bëjmë analiza të mëtejshme, si dhe të ndajmë të dhënat me kolegët dhe të fillojmë të përgjigjemi ndaj pyetjeve të biznesit.
SELECT * FROM `user-logs-237110.userlogs.logdata` LIMIT 10; 
Figura 5: BigQuery
Përfundimi
Shpresojmë që ky postim të shërbejë si një shembull të dobishëm të krijimit të një tubi të dhënash në kohë reale, si dhe gjetjes së mënyrave për ta bërë të dhënat më të arritshme. Ruajtja e të dhënave në një format të tillë na ofron shumë përfitime. Tani mund të fillojmë të përgjigjemi ndaj pyetjeve të rëndësishme, si sa njerëz përdorin produktin tonë? A po rritet baza e përdoruesve me kalimin e kohës? Me cilat aspekte të produktit ndërveprojnë më shumë njerëzit? Dhe a ka gabime, aty ku nuk duhet? Këto janë pyetje që do të jenë të interesuara për organizatën. Bazuar në idetë që del nga përgjigjet në këto pyetje, ne do të jemi në gjendje të përmirësojmë produktin dhe të rrisim angazhimin e përdoruesve.
Beam është vërtet i dobishëm për këtë lloj ushtrimesh, si dhe ka shumë raste të tjera interesante përdorimi. Për shembull, mund të analizoni të dhënat e tregut në kohë reale dhe të bëni tregti të bazuara në analiza, ndoshta keni të dhëna sensorësh që vijnë nga mjete dhe dëshironi të llogaritni nivelin e trafik, gjithashtu, mund të jeni një kompani lojërash që mbledh të dhëna për përdoruesit dhe përdor ato për të krijuar panelet informuese për të ndjekur treguesit kyç. Mirë, shokë, kjo është një temë për një post tjetër, faleminderit për leximin, dhe për ata që duan të shohin kodin e plotë, më poshtë është lidhja për GitHub-in tim.
Këtu mbaron gjithçka. .
Burimi: habr.com
