Loome andmete voogude töötlemise toru. Osa 2

Tere kĂ”igile. Jagame kursuse ĂŒliĂ”pilastele spetsiaalselt valmistatud artikli viimase osa tĂ”lget „Andmeinsener“. Esimese osaga saab tutvuda siin.

Apache Beam ja DataFlow reaalajas konveierite jaoks

Loome andmete voogude töötlemise toru. Osa 2

Google Cloudi seadistamine

MÀrkus: kasutades Google Cloud Shell'i andmete konveieri kÀivitamiseks ja kasutajalogide avaldamiseks, tekkis mul probleeme konveieri kÀivitamisega Python 3 peal. Google Cloud Shell kasutab Python 2, mis sobib paremini Apache Beami jaoks.

Konveieri kÀivitamiseks peame natuke seadistustes tuhnima. Teistest teistest GCP-del pole kasutanud, peate jÀrgima jÀrgmisi 6 sammu, mis on esitatud sellel lehelt.

PÀrast seda peame meie skriptid Google'i pilvesalvestusse laadima ja neid oma Google Cloud Shell'i kopeerima. Laadimine pilvesalvestusse on piisavalt triviaalne (kirjeldus on saadaval siit). Meie failide kopeerimiseks saame avada Google Cloud Shell'i tööriistaribalt, klÔpsates alloleval joonisel 2 vasakul esimesel ikoonil.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 2

Failide kopeerimiseks ja vajalike raamatukogude installimiseks vajalikud kÀsklused 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

Kui oleme lĂ€binud kĂ”ik seadistamise sammud, on jĂ€rgmine asi, mida peame tegema, luua BigQuerys andmestik ja tabel. On mitu viisid, kuidas seda teha, kuid kĂ”ige lihtsam on kasutada Google Cloudi konsooli, alustades andmestiku loomisega. Saate jĂ€rgida jĂ€rgmisi lingi kaudu, et luua tabel oma skeemiga. Meie tabelil on 7 veergu, mis vastavad iga kasutajapĂ”hise logi komponentidele. Mugavuse huvides mÀÀratleme kĂ”ik veerud stringidena (tĂŒĂŒp string), vĂ€lja arvatud muutuja timelocal, ja nimetame need vastavalt varasemalt genereeritud muutujatele. Meie tabeli skeem peaks vĂ€lja nĂ€gema nagu joonisel 3.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 3. Tabeli skeem

KasutajapÔhise logi andmete avaldamine

Pub/Sub on meie voosĂŒsteemi kriitiliselt tĂ€htis komponent, mis vĂ”imaldab mitmel sĂ”ltumatul rakendusel omavahel suhelda. Eriti toimib see vahendajana, vĂ”imaldades meil rakenduste vahel sĂ”numeid saata ja vastu vĂ”tta. Esimene asi, mida peame tegema, on teema (topic) loomine. Piisab, kui minna Pub/Sub-i konsooli ja vajutada LOOMINE.

Allolev kood kutsub meie skripti, et genereerida logi andmeid, nagu varem mÀÀratletud, ning seejĂ€rel ĂŒhendub ja saadab logid Pub/Sub-i. Ainus, mida me peame tegema, on luua objekt PublisherClient, mÀÀrata teema rada meetodiga topic_path ja kutsuda vĂ€lja funktsioon publish koos topic_path ja andmetega. Pange tĂ€hele, et impordime generate_log_line meie skriptist stream_logs, seega veenduge, et need failid on samas kaustas, vastasel juhul saate impordi tĂ”rke. Siis saame seda kĂ€ivitada meie google-konsoolis, kasutades:

python publish.py

import stream_logs
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):
    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 = stream_logs.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 on kÀivitatud, 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.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 4. VĂ€ljund publish_logs.py

Meie vooluahela koodi kirjutamine

NĂŒĂŒd, kui kĂ”ik on valmis, saame alustada kĂ”ige huvitavamat osa — meie konveieri koodi kirjutamist, kasutades Beam'i ja Pythoni. Beam-konveieri loomiseks peame looma konveieri objekti (p). PĂ€rast objekti loomist saame jĂ€rjestikku rakendada mitmeid funktsioone, kasutades operaatorit pipe (|). Üldiselt nĂ€eb töövoog vĂ€lja nagu alloleval joonisel.

[Final Output PCollection] = ([Initial Input PCollection] | [First Transform]
             | [Second Transform]
             | [Third Transform])

Meie koodis loome kaks kohandatud funktsiooni. Funktsiooni regex_clean, mis skaneerib andmeid ja ekstraheerib vastava rea PATTERNS nimekirja pĂ”hjal, kasutades funktsiooni re.search. Funktsioon tagastab koma kaudu eraldatud rea. Kui te pole regulaarsete vĂ€ljendite spetsialist, soovitan tutvuda selle Ă”pikuga ja harjutada mĂ€rkmikus, et koodi katsetada. PĂ€rast seda mÀÀratleme kohandatud ParDo-funktsiooni nimega Split, mis on Beam-transformatsiooni variatsioon paralleelse töötlemise jaoks. Pythonis tehakse seda spetsiifilisel viisil — me peame looma klassi, mis pĂ€rib Beam’i DoFn klassist. Funktsioon Split vĂ”tab eelmisest funktsioonist parsetud stringi ja tagastab sĂ”nastike loendi, mille vĂ”tmed vastavad meie BigQuery tabeli veergude nimedele. On mĂ”ned asjad, mida tuleb selle funktsiooni kohta mĂ€rkida: pidin selle tööle saamiseks impordima datetime sisemiselt funktsiooni. Saime alguses faili impordimisel veateate, mis oli kummaline. See loend antakse seejĂ€rel funktsioonile WriteToBigQuery, mis lihtsalt lisab meie andmed tabelisse. Koodi Batch DataFlow Job'i ja Streaming DataFlow Job'i jaoks on esitatud allpool. Ainus erinevus partii ja voogesituse koodi vahel on see, et partii töötlemisel loeme CSV-d src_path, kasutades funktsiooni ReadFromText Beam'ist.

Batch DataFlow Job (partiis 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

PROJEKT='user-logs-237110'
skeem = '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):

    MUSTRID =  [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
           r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(? beam.io.textio.ReadFromText(src_path)
      | "puhasta aadress" > beam.Map(regex_clean)
      | 'ParseCSV' > beam.ParDo(Split())
      | 'WriteToBigQuery' > beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJEKT), schema=skeem,
        write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
   )

   p.run()

if __name__ == '__main__':
  logger = logging.getLogger().setLevel(logging.INFO)
  main()

Streaming DataFlow Job (voogesituse töötlemine)

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)
      | "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()

Konveieri kÀivitamine

Saame kÀivitada konveieri mitmel erineval viisil. Kui soovime, saame selle lihtsalt kÀivitada kohalikult terminalist, sisenedes eemalt GCP-sse.

python -m main_pipeline_stream.py 
 --input_topic "projects/user-logs-237110/topics/userlogs" 
 --streaming

Kuid me kavatseme selle kÀivitada DataFlow abil. Saame seda teha alloleva kÀsu abil, seadistades jÀrgmised kohustuslikud parameetrid.

  • project — teie GCP projekti ID.
  • runner — konveieri kĂ€ivitamise vahend, mis analĂŒĂŒsib teie programmi ja konstruerib teie konveieri. Pilves töötamiseks peate mÀÀrama DataflowRunner.
  • staging_location — rada Cloud Dataflow pilve salvestusse, et indekseerida koodipakette, mida vajavad töötajad, kes tööd teevad.
  • temp_location — rada Cloud Dataflow pilve salvestusse ajutiste failide asukoha mÀÀramiseks, mis loodakse konveieri töö ajal.
  • streaming

python main_pipeline_stream.py 
--runner DataFlow 
--project $PROJECT 
--temp_location $BUCKET/tmp 
--staging_location $BUCKET/staging
--streaming

Kuna see kÀsk on töös, saame minna google'i konsooli DataFlow vahekaardile ja vaadata meie toru. Torule klikkides peaksime nÀgema midagi sarnast joonisele 4. TÔrkeotsingu vajadusel vÔib olla vÀga kasulik minna logidesse ja seejÀrel Stackdriverisse detailsete logide vaatamiseks. See on aidanud mul mitmel korral toru probleeme lahendada.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 4: Beam-toru

Juura meie andmetele BigQuery's

Seega peaks meil juba olema tööle pandud toru, kuhu andmed voolavad meie tabelisse. Selle kontrollimiseks saame minna BigQuery'sse ja vaadata andmeid. PĂ€rast alltoodud kĂ€skluse kasutamist peaksite nĂ€gema andmekogumi esimesi ridasid. NĂŒĂŒd, kui meil on andmed hoitud BigQuery's, saame teha edasist analĂŒĂŒsi, jagada andmeid kolleegidega ja hakata vastama Ă€rikĂŒsimustele.

SELECT * FROM `user-logs-237110.userlogs.logdata` LIMIT 10;

Loome andmete voogude töötlemise toru. Osa 2
Joonis 5: BigQuery

KokkuvÔte

Loodame, et see postitus annab kasuliku nĂ€ite andmevoogude loomisest ning toob esile viise, kuidas andmeid kergemini kĂ€tte saada. Andmete selle vormingus hoidmine pakub meile mitmeid eeliseid. NĂŒĂŒd saame hakata vastama olulistele kĂŒsimustele, nagu nĂ€iteks, kui palju inimesi meie toodet kasutab? Kas kasutajate arv kasvab ajas? Milliste toote aspektidega inimesed kĂ”ige rohkem suhtlevad? Ja kas on vigu seal, kus neid ei peaks olema? Need on kĂŒsimused, mis organisatsiooni jaoks olulised. Vastuste pĂ”hjal, mis tulenevad nendest kĂŒsimustest, saame toote tĂ€iustamiseks ja kasutajate huvi suurendamiseks hĂ€id ideid.

Beam on tĂ”eliselt kasulik selliste harjutuste jaoks ning sellel on palju muid huvitavaid kasutusjuhtumeid. NĂ€iteks vĂ”ite analĂŒĂŒsida börsi tehingute andmeid reaalajas ja teha tehinguid analĂŒĂŒsi pĂ”hjal. VĂ”ib-olla on teil andurid, mis saadavad andmeid sĂ”idukitest, ja soovite arvutada liiklusvoogu. Samuti vĂ”ite olla mĂ€nguettevĂ”te, kes kogub kasutajaandmeid ning kasutab neid ĂŒlevaatetablettide loomisel, et jĂ€lgida vĂ”tmemÔÔdikute tĂ€itmist. Noh, hĂ€rrased, see on juba teine teema, tĂ€nan lugemise eest, ning neile, kes soovivad nĂ€ha tĂ€ielikku koodi, leiate allpool lingi minu GitHubi lehele.

https://github.com/DFoly/User_log_pipeline

Sellega on kÔik. Loe esimest osa.

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster