Création d'un pipeline de traitement de données. Partie 2

Bonjour à tous. Nous partageons la traduction de la dernière partie d'un article préparé spécialement pour les étudiants du cours «Ingénieur en données». Vous pouvez consulter la première partie ici.

Apache Beam et DataFlow pour des pipelines en temps réel

Création d'un pipeline de traitement de données. Partie 2

Configuration de Google Cloud

Remarque : Pour exécuter le pipeline et publier les données du journal personnalisé, j'ai utilisé Google Cloud Shell, car j'ai rencontré des problèmes pour exécuter le pipeline sur Python 3. Google Cloud Shell utilise Python 2, qui est mieux compatible avec Apache Beam.

Pour exécuter le pipeline, nous devons creuser un peu dans les paramètres. Pour ceux d'entre vous qui n'ont jamais utilisé GCP, vous devrez suivre les 6 étapes décrites sur ce la page.

Après cela, nous devrons télécharger nos scripts dans le stockage cloud de Google et les copier dans notre Google Cloud Shell. Le téléchargement dans le stockage cloud est assez trivial (vous pouvez trouver la description ici). Pour copier nos fichiers, nous pouvons ouvrir Google Cloud Shell à partir de la barre d'outils en cliquant sur la première icône à gauche sur l'image 2 ci-dessous.

Création d'un pipeline de traitement de données. Partie 2
Figure 2

Les commandes dont nous avons besoin pour copier des fichiers et installer les bibliothèques nécessaires sont énumérées ci-dessous.

# 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>

Création de notre base de données et table

Une fois que nous avons suivi toutes les étapes de configuration, la prochaine chose à faire est de créer un ensemble de données et une table dans BigQuery. Il existe plusieurs façons de le faire, mais la plus simple consiste à utiliser la console Google Cloud, en créant d'abord un ensemble de données. Vous pouvez suivre les étapes indiquées ci-dessous le lien, pour créer une table avec le schéma. Notre table aura 7 colonnes, correspondant aux composants de chaque journal personnalisé. Pour simplifier, nous définirons toutes les colonnes comme des chaînes (type string), sauf la variable timelocal, et nous les nommerons en fonction des variables que nous avons générées précédemment. Le schéma de notre table devrait ressembler à l'image 3.

Création d'un pipeline de traitement de données. Partie 2
Image 3. Schéma de la table

Publication des données du journal personnalisé

Pub/Sub est un composant critique de notre pipeline, car il permet à plusieurs applications indépendantes d'interagir les unes avec les autres. En particulier, il fonctionne comme un intermédiaire, nous permettant d'envoyer et de recevoir des messages entre les applications. La première chose à faire est de créer un sujet. Il suffit d'aller dans Pub/Sub dans la console et de cliquer sur CREER UN SUJET.

Le code ci-dessous appelle notre script pour générer les données de journal comme décrit ci-dessus, puis se connecte et envoie les journaux à Pub/Sub. La seule chose que nous devons faire est de créer un objet PublisherClient, spécifier le chemin du sujet via la méthode topic_path et appeler la fonction publish avec topic_path avec les données. Notez que nous importons generate_log_line de notre script stream_logs, donc assurez-vous que ces fichiers se trouvent dans le même dossier, sinon vous obtiendrez une erreur d'importation. Ensuite, nous pouvons exécuter cela via notre console Google, en utilisant :

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):
    # Lorsque le délai d'attente n'est pas spécifié, la méthode d'exception attend indéfiniment.
    if message_future.exception(timeout=30):
        print('La publication du message sur {} a provoqué une 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)

Une fois que le fichier sera exécuté, nous pourrons observer la sortie des données de journal sur la console, comme montré sur l'image ci-dessous. Ce script continuera à fonctionner tant que nous n'utiliserons pas CTRL+C, pour l'arrêter.

Création d'un pipeline de traitement de données. Partie 2
Figure 4. Sortie publish_logs.py

Écriture du code de notre pipeline

Maintenant que nous avons tout préparé, nous pouvons passer à la partie la plus intéressante — écrire le code de notre pipeline en utilisant Beam et Python. Pour créer un pipeline Beam, nous devons créer un objet pipeline (p). Une fois que nous avons créé l'objet pipeline, nous pouvons appliquer plusieurs fonctions les unes après les autres, en utilisant l'opérateur pipe (|). En général, le flux de travail ressemble à l'image ci-dessous.

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

Dans notre code, nous allons créer deux fonctions personnalisées. La fonction regex_clean, qui analyse les données et extrait la chaîne correspondante basée sur la liste PATTERNS, en utilisant la fonction re.search. La fonction retourne une chaîne séparée par des virgules. Si vous n'êtes pas un expert en expressions régulières, je vous conseille de consulter ce tutoriel et de pratiquer dans le bloc-notes pour vérifier le code. Après cela, nous définissons une fonction ParDo personnalisée appelée Split, qui est une variation de la transformation Beam pour le traitement parallèle. En Python, cela se fait d'une manière particulière : nous devons créer une classe qui hérite de la classe DoFn de Beam. La fonction Split prend une chaîne analysée de la fonction précédente et renvoie une liste de dictionnaires avec des clés correspondant aux noms des colonnes de notre table BigQuery. Il y a quelque chose à noter à propos de cette fonction : j'ai dû importer datetime à l'intérieur de la fonction pour qu'elle fonctionne. J'ai reçu un message d'erreur lors de l'importation en début de fichier, ce qui était étrange. Cette liste est ensuite transmise à la fonction WriteToBigQuery, qui ajoute simplement nos données dans la table. Le code pour le Batch DataFlow Job et le Streaming DataFlow Job est donné ci-dessous. La seule différence entre le code par lots et le code en streaming est que dans le traitement en lots, nous lisons un CSV à partir de src_path, en utilisant la fonction ReadFromText de Beam.

Batch DataFlow Job (traitement par lots)

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)
      | "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 (traitement en continu)

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

Démarrer le pipeline

Nous pouvons lancer le pipeline de plusieurs manières. Si nous le souhaitions, nous pourrions simplement l'exécuter localement depuis le terminal, en nous connectant à distance à GCP.

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

Cependant, nous allons le lancer en utilisant DataFlow. Nous pouvons le faire avec la commande ci-dessous en spécifiant les paramètres obligatoires suivants.

  • project — L'ID de votre projet GCP.
  • runner — Le moteur d'exécution du pipeline, qui analysera votre programme et construira votre pipeline. Pour l'exécution dans le cloud, vous devez spécifier DataflowRunner.
  • staging_location — Le chemin vers le stockage cloud de Cloud Dataflow pour l'indexation des paquets de code nécessaires aux workers exécutant le travail.
  • temp_location — Le chemin vers le stockage cloud de Cloud Dataflow pour le stockage des fichiers temporaires des tâches créés pendant l'exécution du pipeline.
  • streaming

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

Pendant que cette commande s'exécute, nous pouvons passer à l'onglet DataFlow dans la console Google et consulter notre pipeline. En cliquant sur le pipeline, nous devrions voir quelque chose comme l'image 4. Pour le débogage, il peut être très utile de consulter les journaux, puis Stackdriver pour voir les journaux détaillés. Cela m'a aidé à résoudre des problèmes de pipeline à plusieurs reprises.

Création d'un pipeline de traitement de données. Partie 2
Image 4 : Pipeline Beam

Accéder à nos données dans BigQuery

Donc, nous devrions déjà avoir un pipeline en cours d'exécution avec des données entrant dans notre tableau. Pour vérifier cela, nous pouvons nous rendre dans BigQuery et consulter les données. Après avoir utilisé la commande ci-dessous, vous devriez voir les premières lignes de l'ensemble de données. Maintenant que nous avons des données stockées dans BigQuery, nous pouvons procéder à une analyse plus approfondie, ainsi que partager les données avec des collègues et commencer à répondre aux questions commerciales.

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

Création d'un pipeline de traitement de données. Partie 2
Image 5 : BigQuery

Conclusion

Nous espérons que ce post servira d'exemple utile pour créer un pipeline de données en streaming, ainsi que pour trouver des moyens de rendre les données plus accessibles. Le stockage des données dans ce format nous offre de nombreux avantages. Nous pouvons maintenant commencer à répondre à des questions importantes, telles que : combien de personnes utilisent notre produit ? La base d'utilisateurs augmente-t-elle avec le temps ? Quelles fonctionnalités du produit interagissent le plus ? Et y a-t-il des erreurs là où elles ne devraient pas être ? Ce sont des questions qui intéresseront l'organisation. Sur la base des idées résultant des réponses à ces questions, nous pourrons améliorer le produit et accroître l'engagement des utilisateurs.

Beam est vraiment utile pour ce type d'exercices, et il a également un certain nombre d'autres cas d'utilisation intéressants. Par exemple, vous pouvez analyser des données de ticks boursiers en temps réel et prendre des décisions commerciales basées sur cette analyse, peut-être que vous disposez de données de capteurs provenant de véhicules et que vous souhaitez calculer le niveau de trafic. Vous pouvez également, par exemple, être une entreprise de jeux collectant des données sur les utilisateurs et les utilisant pour créer des tableaux de bord afin de suivre les indicateurs clés. Bien, messieurs, c'est un sujet pour un autre пост, merci de votre lecture, et pour ceux qui souhaitent voir le code complet, voici le lien vers mon GitHub.

https://github.com/DFoly/User_log_pipeline

C'est tout. Lire la première partie.

Source : habr.com

Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS 🔥 Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster