Hola a todos. Compartimos la traducción de la parte final de un artículo preparado especialmente para los estudiantes del curso . Puedes consultar la primera parte .
Apache Beam y DataFlow para tuberías de tiempo real

Configuración de Google Cloud
Nota: Para ejecutar la tubería y publicar los datos del registro del usuario, utilicé Google Cloud Shell, ya que tuve problemas para ejecutar la tubería en Python 3. Google Cloud Shell utiliza Python 2, que es más compatible con Apache Beam.
Para ejecutar la tubería, necesitamos profundizar un poco en la configuración. Aquellos de ustedes que no han utilizado GCP antes deberán seguir estos 6 pasos que se indican en esta .
Después de esto, necesitaremos cargar nuestros scripts en el almacenamiento en la nube de Google y copiarlos a nuestro Google Cloud Shell. La carga en el almacenamiento en la nube es bastante trivial (la descripción se puede encontrar ). Para copiar nuestros archivos, podemos abrir Google Cloud Shell desde la barra de herramientas, haciendo clic en el primer ícono a la izquierda en la figura 2 a continuación.

Figura 2
Los comandos que necesitamos para copiar archivos e instalar las bibliotecas necesarias se enumeran a continuación.
# 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>Creación de nuestra base de datos y tabla
Una vez que hayamos completado todos los pasos relacionados con la configuración, lo siguiente que debemos hacer es crear un conjunto de datos y una tabla en BigQuery. Hay varias maneras de hacerlo, pero la más sencilla es usar la consola de Google Cloud, creando primero un conjunto de datos. Puedes seguir los pasos que se indican a continuación , para crear una tabla con el esquema. Nuestra tabla tendrá 7 columnas, que corresponden a los componentes de cada registro de usuario. Para simplificar, definiremos todas las columnas como cadenas (tipo string), excepto la variable timelocal, y las nombraremos de acuerdo con las variables que generamos anteriormente. El esquema de nuestra tabla debería verse como se muestra en la figura 3.

Figura 3. Esquema de la tabla
Publicación de datos del registro del usuario
Pub/Sub es un componente crítico de nuestra tubería, ya que permite a varias aplicaciones independientes interactuar entre sí. En particular, funciona como un intermediario que nos permite enviar y recibir mensajes entre aplicaciones. Lo primero que debemos hacer es crear un tema (topic). Es bastante sencillo ir a Pub/Sub en la consola y hacer clic en CREAR TEMA.
El código a continuación llama a nuestro script para generar datos de registro definidos anteriormente y luego se conecta y envía los registros a Pub/Sub. Lo único que necesitamos hacer es crear un objeto PublisherClient, especificar la ruta del tema utilizando el método topic_path y llamar a la función publish con topic_path y los datos. Tenga en cuenta que estamos importando generate_log_line de nuestro script stream_logs, así que asegúrese de que estos archivos estén en la misma carpeta, de lo contrario obtendrá un error de importación. Luego podemos ejecutar esto a través de nuestra consola de Google usando:
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):
# Cuando el tiempo de espera no está especificado, el método de excepción espera indefinidamente.
if message_future.exception(timeout=30):
print('Publicar el mensaje en {} lanzó una excepción {}.'.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)Una vez que se ejecute el archivo, podremos observar la salida de los datos de registro en la consola, como se muestra en la imagen a continuación. Este script funcionará hasta que usemos CTRL+C, para finalizarlo.

Figura 4. Salida publish_logs.py
Escribiendo el código de nuestro pipeline
Ahora que hemos preparado todo, podemos pasar a la parte más interesante: escribir el código de nuestro pipeline usando Beam y Python. Para crear un pipeline de Beam, necesitamos crear un objeto de pipeline (p). Una vez que hemos creado el objeto de pipeline, podemos aplicar varias funciones una tras otra utilizando el operador pipe (|). En general, el flujo de trabajo se ve como en la imagen a continuación.
[Final Output PCollection] = ([Initial Input PCollection] | [First Transform]
| [Second Transform]
| [Third Transform]) En nuestro código, crearemos dos funciones personalizadas. La función regex_clean, que escanea los datos y extrae la cadena correspondiente basada en la lista de PATTERNS utilizando la función re.search. La función devuelve una cadena separada por comas. Si no es un experto en expresiones regulares, le recomiendo que consulte este y practicar en un cuaderno para comprobar el código. Después de esto, definimos una función ParDo personalizada llamada Split, que es una variación de la transformación Beam para el procesamiento paralelo. En Python, esto se hace de una manera particular: debemos crear una clase que herede de la clase DoFn de Beam. La función Split toma una cadena analizada de la función anterior y devuelve una lista de diccionarios con claves que corresponden a los nombres de las columnas en nuestra tabla de BigQuery. Hay algo que vale la pena mencionar sobre esta función: tuve que importar datetime dentro de la función para que funcionara. Recibía un mensaje de error al importar al principio del archivo, lo que era extraño. Esta lista se pasa luego a la función WriteToBigQuery, que simplemente añade nuestros datos a la tabla. El código para Batch DataFlow Job y Streaming DataFlow Job se presenta a continuación. La única diferencia entre el código por lotes y el código en streaming es que en el procesamiento por lotes leemos CSV de src_path, usando la función ReadFromText de Beam.
Batch DataFlow Job (procesamiento por lotes)
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("Hubo un error con la búsqueda regex")
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 (procesamiento de flujo)
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()
Ejecutar el canal
Podemos ejecutar el canal de varias maneras diferentes. Si quisiéramos, podríamos simplemente ejecutarlo localmente desde la terminal, ingresando de forma remota en GCP.
python -m main_pipeline_stream.py
--input_topic "projects/user-logs-237110/topics/userlogs"
--streamingSin embargo, vamos a ejecutarlo utilizando DataFlow. Podemos hacerlo con el siguiente comando, estableciendo los parámetros obligatorios a continuación.
project— ID de su proyecto GCP.runner— medio para ejecutar el canal que analizará su programa y construirá su canal. Para la ejecución en la nube, debe especificar DataflowRunner.staging_location— ruta al almacenamiento en la nube de Cloud Dataflow para indexar los paquetes de código necesarios para los manejadores que ejecutan el trabajo.temp_location— ruta al almacenamiento en la nube de Cloud Dataflow para ubicar los archivos temporales de las tareas creadas durante la ejecución del canal.streaming
python main_pipeline_stream.py
--runner DataFlow
--project $PROJECT
--temp_location $BUCKET/tmp
--staging_location $BUCKET/staging
--streaming
Mientras se ejecuta este comando, podemos ir a la pestaña DataFlow en la consola de Google y revisar nuestro pipeline. Al hacer clic en el pipeline, deberíamos ver algo similar a la figura 4. Para fines de depuración, puede ser muy útil ir a los registros y luego a Stackdriver para ver los registros detallados. Esto me ha ayudado a resolver problemas con el pipeline en varias ocasiones.

Figura 4: Pipeline de Beam
Acceso a nuestros datos en BigQuery
Así que ya deberíamos tener un pipeline en ejecución con datos fluyendo hacia nuestra tabla. Para verificar esto, podemos ir a BigQuery y ver los datos. Después de usar el comando a continuación, deberías ver las primeras filas del conjunto de datos. Ahora que tenemos datos almacenados en BigQuery, podemos realizar un análisis adicional, así como compartir los datos con colegas y comenzar a responder preguntas comerciales.
SELECT * FROM `user-logs-237110.userlogs.logdata` LIMIT 10; 
Figura 5: BigQuery
Conclusión
Esperamos que esta publicación sirva como un ejemplo útil para la creación de un pipeline de datos en flujo, así como para encontrar formas de hacer los datos más accesibles. Almacenar datos en este formato nos brinda muchas ventajas. Ahora podemos comenzar a responder preguntas importantes, como ¿cuántas personas utilizan nuestro producto? ¿Está creciendo nuestra base de usuarios con el tiempo? ¿Con qué aspectos del producto interactúan más las personas? ¿Y hay errores donde no deberían estar? Estas son preguntas que serán de interés para la organización. Basándonos en las ideas que surgen de las respuestas a estas preguntas, podremos mejorar el producto y aumentar el interés de los usuarios.
Beam es realmente útil para este tipo de ejercicios, y también tiene una serie de otros casos de uso interesantes. Por ejemplo, puedes analizar datos de tickers de bolsa en tiempo real y realizar transacciones basadas en ese análisis. Tal vez tengas datos de sensores que provienen de vehículos y quieras calcular el nivel de tráfico. También podrías ser una empresa de videojuegos recopilando datos de usuarios y utilizándolos para crear tableros de control para rastrear indicadores clave. Bueno, señores, este es un tema para otra publicación, gracias por leer, y para aquellos que desean ver el código completo, aquí abajo está el enlace a mi GitHub.
Eso es todo. .
Fuente: habr.com
