Tere kÔigile. Jagame tÔlget artikli lÔpposa, mis on spetsiaalselt ette valmistatud kursuse tudengitele. . Esimese osaga saab tutvuda .
Apache Beam ja DataFlow reaalaja töötluse jaoks

Google Cloudi seadistamine
MÀrkus: Konveieri kÀivitamiseks ja kohandatud logi andmete avaldamiseks kasutasin Google Cloud Shelli, kuna mul oli probleem konveieri kÀivitamisega Python 3-s. Google Cloud Shell kasutab Python 2-d, mis sobib paremini Apache Beamiga.
Konveieri kÀivitamiseks peame natuke seadistuses kaevuma. Neil teist, kes pole varem GCP-d kasutanud, tuleb jÀrgida jÀrgmisi 6 sammu, nagu on esitatud .
PÀrast seda tuleb meil laadida meie skriptid Google'i pilve salvestusse ja kopeerida need meie Google Cloud Shel'i. Laadimine pilve salvestusse on piisavalt lihtne (kirjeldus on saadaval ). Meie failide kopeerimiseks saame avada Google Cloud Shel'i tööriistaribalt, klÔpsates alloleval pildil 2 esimesel ikoonil.

Joonis 2
Failide kopeerimiseks ja vajalike teekide installimiseks vajalikud kÀsud on loetletud allpool.
# 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>Meie andmebaasi ja tabeli loomine
PĂ€rast kĂ”ikide seadistamisega seotud sammude tĂ€itmist on jĂ€rgmine asi, mida teha, luua BigQuery's andmekogum ja tabel. Selle tegemiseks on mitmeid viise, kuid kĂ”ige lihtsam on kasutada Google Cloudi konsoli, luues esmalt andmekogumi. Saate jĂ€rgida jĂ€rgmisi , et luua tabel skeemiga. Meie tabelil on 7 veergu, mis vastavad iga kohandatud logi komponentidele. Mugavuse huvides mÀÀratleme kĂ”ik veerud string-tĂŒĂŒpi (string), vĂ€lja arvatud muutuja timelocal, ja nimetame need vastavalt meie varem genereeritud muutujatele. Meie tabeli skeem peaks vĂ€lja nĂ€gema nagu joonisel 3.

Joonis 3. Tabeli skeem
Kohandatud logi andmete avaldamine
Pub/Sub on meie konveieri kriitiliselt oluline komponent, kuna see vÔimaldab mitmetel sÔltumatutel rakendustel omavahel suhelda. EelkÔige toimib see vahendajana, mis vÔimaldab meil edastada ja vastu vÔtta sÔnumeid rakenduste vahel. Esimene asi, mida me peame tegema, on luua teema (topic). Piisab lihtsalt Pub/Sub-i minemisest konsoolis ja vajutades LOO TEEMA.
Allpoolt esitatud kood kutsub vĂ€lja meie skripti logiandmete genereerimiseks, nagu eespool mÀÀratletud, ja seejĂ€rel ĂŒhendub ja saadab logid Pub/Subi. Ainus, mida peame tegema, on luua objekt PublisherClient, mÀÀrata teema rada meetodi kaudu topic_path ja kutsuda vĂ€lja funktsioon publish jot topic_path ja andmete poolt. Pange tĂ€hele, et impordime generate_log_line meie skriptist stream_logs, seega veenduge, et need failid asuksid samas kaustas, muidu saate importimisvea. SeejĂ€rel saame selle kĂ€ivitada meie google'i konsoolist, kasutades:
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):
# Kui tÀhtaeg pole mÀÀratud, siis ootab erandi meetod lÔputult.
if message_future.exception(timeout=30):
print('SÔnumi avaldamine {} tÔi kaasa erandi {}.'.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)Kui fail kÀivitub, saame jÀlgida logiandmete vÀljundit konsoolis, nagu on nÀidatud alloleval joonisel. See skript töötab seni, kuni me ei kasuta CTRL+C, et see lÔpetada.

Joonis 4. VĂ€ljund publish_logs.py
Meie toru koodi kirjutamine
NĂŒĂŒd, kui oleme kĂ”ik ette valmistanud, saame liikuda kĂ”ige huvitavama osa juurde - meie toru koodi kirjutamine, kasutades Beam'i ja Pythonit. Beam toru loomiseks peame looma toru objekti (p). PĂ€rast selle objekti loomist saame rakendada mitmeid funktsioone jĂ€rjestikku, kasutades operaatorit pipe (|). Ăldiselt nĂ€eb töövoog vĂ€lja nagu alloleval joonisel.
[LÔplik VÀljund PCollection] = ([Algne Sisend PCollection] | [Esimene Muundamine]
| [Teine Muundamine]
| [Kolmas Muundamine]) Meie koodis loome kaks kasutaja mÀÀratud funktsiooni. Funktsiooni regex_clean, mis skaneerib andmed ja ekstraktib vastava rea, tuginedes PATTERNS nimekirjale, kasutades funktsiooni re.search. Funktsioon tagastab koma jagatud rea. Kui te ei ole regulaarsete vĂ€ljendite ekspert, soovitan tutvuda selle ja harjutada mĂ€rkmikus koodi testimiseks. PĂ€rast seda mÀÀratleme kasutaja ParDo-funktsiooni nimega Split, mis on Beam'i transformatsiooni variant paralleelse töötlemise jaoks. Pythonis tehakse seda erilisel viisil â peame looma klassi, mis pĂ€rib Beam'i klassist DoFn. Funktsioon Split vĂ”tab eelmisest funktsioonist lahendatud stringi ja tagastab sĂ”nastike loendi, mille vĂ”tmed vastavad meie BigQuery tabeli veergude nimedele. On siiski midagi, mida tuleb selle funktsiooni kohta mĂ€rkida: pidin impordima datetime funktsiooni sees, et see töötaks. Sain alguses faili impordimisel veateate, mis oli kummaline. See loend edastatakse seejĂ€rel funktsioonile WriteToBigQuery, mis lihtsalt lisab meie andmed tabelisse. Batch DataFlow Job ja Streaming DataFlow Job kood on toodud allpool. Ainuke erinevus partii ja voolu koodi vahel on see, et partii töötlemisel loeme CSV-d aadressilt src_path, kasutades funktsiooni ReadFromText Beam'ist.
Batch DataFlow Job (pakettide töötlemine)
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"(?<=[).d+(?=])",
r'"[A-Z][a-z]+', r'"(http|https):\/\/[a-z]+.[a-z]+.[a-z]+']
result = []
for match in PATTERNS:
try:
reg_match = re.search(match, data).group()
if reg_match:
result.append(reg_match)
else:
result.append(" ")
except:
print("Regex otsingus tekkis viga")
result = [x.strip() for x in result]
result = [x.replace('"', "") for x in result]
res = ','.join(result)
return res
class Split(beam.DoFn):
def process(self, element):
from datetime import datetime
element = element.split(",")
d = datetime.strptime(element[1], "%d\/ %b\/ %Y:%H:%M:%S")
date_string = d.strftime("%Y-%m-%d %H:%M:%S")
return [{
'remote_addr': element[0],
'timelocal': date_string,
'request_type': element[2],
'status': element[3],
'body_bytes_sent': element[4],
'http_referer': element[5],
'http_user_agent': element[6]
}]
def main():
p = beam.Pipeline(options=PipelineOptions())
(p
| 'ReadData' >> beam.io.textio.ReadFromText(src_path)
| "clean address" >> beam.Map(regex_clean)
| 'ParseCSV' >> beam.ParDo(Split())
| 'WriteToBigQuery' >> 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 (voo töötlemine)
import apache_beam.options.pipeline_options as 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)
| "Decode" >> beam.Map(lambda x: x.decode('utf-8'))
| "Clean Data" >> beam.Map(regex_clean)
| 'ParseCSV' >> beam.ParDo(Split())
| 'WriteToBigQuery' >> 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()
Konteineri kÀivitamine
Saame vÔime kÀivitada torujuhtme mitmel erineval viisil. Kui soovime, saame selle lihtsalt kohapeal terminali kaudu kÀivitada, sisenedes kaugjuhtimisega GCP-sse.
python -m main_pipeline_stream.py
--input_topic "projects/user-logs-237110/topics/userlogs"
--streamingKuid me kavatseme selle kÀivitada DataFlow kaudu. Saame seda teha allpool toodud kÀsuga, seadistades jÀrgmised kohustuslikud parameetrid.
projectâ Teie GCP projekti ID.runnerâ torujuhtme kĂ€ivitamise tööriist, mis analĂŒĂŒsib teie programmi ja konstrueerib teie torujuhtme. Pilve tĂ€itmiseks peate mÀÀrama DataflowRunneri.staging_locationâ tee Cloud Dataflow pilvealusesse salvestusse koodipakettide indekseerimiseks, mis on vajalik tööde tĂ€itmiseks.temp_locationâ tee Cloud Dataflow pilvealusesse salvestusse ajutiste failide jaoks, mis on loodud torujuhtme kĂ€itamise ajal.voog
python main_pipeline_stream.py
--runner DataFlow
--project $PROJECT
--temp_location $BUCKET/tmp
--staging_location $BUCKET/staging
--streaming
Kuna selle kĂ€skluse tĂ€itmine kestab, saame minna Google'i konsooli vahekaardile DataFlow ja vaadata meie toru. Kui klikime torule, peaksime nĂ€gema midagi, mis sarnaneb joonisele 4. Veaparanduse jaoks vĂ”ib olla vĂ€ga kasulik vaadata logisid ja seejĂ€rel Stackdriveri ĂŒksikasjalikke logisid. See on aidanud mul mitmeid probleeme toruga lahendada.

Joonis 4: Beam-toru
LigipÀÀs meie andmetele BigQuery's
Seega peaks meil juba olema töötav toru, kuhu andmed edastatakse meie tabelisse. Selle kontrollimiseks saame minna BigQuery'sse ja vaadata andmeid. PĂ€rast alloleva kĂ€su kasutamist peaksid ilmuma esimesed paar rida andmekogumist. NĂŒĂŒd, kui meil on andmed BigQuery's, saame teha edasist analĂŒĂŒsi, jagada andmeid kolleegidega ja alustada Ă€rikĂŒsimustele vastamist.
SELECT * FROM `user-logs-237110.userlogs.logdata` LIMIT 10; 
Joonis 5: BigQuery
KokkuvÔte
Loodame, et see postitus toimib kasuliku nĂ€itena voogedastustoru loomise osas ning annab ideid, kuidas andmeid kergemini kĂ€tte saada. Andmete sĂ€ilitamine sellises vormingus pakub meile palju eeliseid. NĂŒĂŒd saame hakata vastama olulistele kĂŒsimustele, nĂ€iteks: kui palju inimesi kasutab meie toodet? Kas kasutajate arv kasvab ajaga? Milliste tootekomponentidega inimesed kĂ”ige rohkem suhtlevad? Ja kas on vigu seal, kus neid olema ei peaks? Need on kĂŒsimused, mis on organisatsiooni jaoks huvitavad. Vastuste pĂ”hjal nendest kĂŒsimustest saame toote paremaks muuta ja kasutajate huvi suurendada.
Beam on tĂ”epoolest kasulik selliste harjutuste jaoks ning tal on ka mitu muud huvitavat kasutusala. NĂ€iteks vĂ”ite analĂŒĂŒsida börsi tehingute andmeid reaalajas ning teha tehinguid analĂŒĂŒsi pĂ”hjal. VĂ”imalik, et teil on sensoriandmed, mis pĂ€rinevad sĂ”idukitest, ja soovite arvutada liiklustaseme. Samuti vĂ”ite olla mĂ€nguettevĂ”te, mis kogub kasutajaandmeid ja kasutab neid teabepaneelide loomiseks vĂ”tmeindikaatorite jĂ€lgimiseks. Noh, sĂ”brad, see on juba teema jĂ€rgmise postituse jaoks, aitĂ€h lugemise eest, ja nende jaoks, kes soovivad tĂ€ielikku koodi nĂ€ha, on allpool link minu GitHubi lehele.
Sellega on kÔik. .
Allikas: habr.com
