Loodame andmevoogude töötlemise toru. Osa 2

Tere kĂ”igile. Jagame tĂ”lget artikli lĂ”pposa, mis on spetsiaalselt ette valmistatud kursuse tudengitele. „Andmeinsener“. Esimese osaga saab tutvuda siit.

Apache Beam ja DataFlow reaalaja töötluse jaoks

Loodame andmevoogude töötlemise toru. Osa 2

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 lehe.

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 siin). Meie failide kopeerimiseks saame avada Google Cloud Shel'i tööriistaribalt, klÔpsates alloleval pildil 2 esimesel ikoonil.

Loodame andmevoogude töötlemise toru. Osa 2
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 lingil, 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.

Loodame andmevoogude töötlemise toru. Osa 2
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.py

from 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.

Loodame andmevoogude töötlemise toru. Osa 2
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 tutoriaga 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" 
 --streaming

Kuid 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.

Loodame andmevoogude töötlemise toru. Osa 2
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;

Loodame andmevoogude töötlemise toru. Osa 2
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.

https://github.com/DFoly/User_log_pipeline

Sellega on kÔik. Loe esimest osa.

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster