Lab práctico · Semana 6: Procesamiento de datos y analítica
Ejecutar y observar un job de Dataflow con una plantilla de Google
Qué vas a construir
Vas a lanzar la plantilla WordCount que Google publica para Dataflow: un job por lotes que lee El rey Lear de un bucket público, cuenta las palabras y escribe el resultado en tu bucket. No escribes código de Apache Beam; el objetivo es ver cómo Dataflow crea workers, ejecuta el grafo, autoescala y termina, y aprender a leer un job en la consola.
flowchart LR
T["Plantilla Word_Count (gs://dataflow-templates)"] --> J["Job de Dataflow en europe-southwest1"]
IN["gs://dataflow-samples/shakespeare/kinglear.txt"] --> J
J --> W["Workers (VMs de Compute Engine gestionadas por Dataflow)"]
J --> OUT["Tu bucket: results/output-*"]
SA["Cuenta de servicio de workers (mínimo privilegio)"] -.-> W
Antes de empezar
- Proyecto dedicado con facturación y Cloud Shell abierto.
- Los workers de Dataflow se crean en la red
default. Si la borraste en el lab 6, añade al comando del paso 3--subnetwork=regions/europe-southwest1/subnetworks/TU_SUBRED(y, si los workers no tienen IP pública, la subred necesita Private Google Access).
export PROJECT_ID=$(gcloud config get-value project)
export REGION=europe-southwest1
export BUCKET=gs://${PROJECT_ID}-lab13
export SA=dataflow-worker-lab13
export SA_EMAIL=${SA}@${PROJECT_ID}.iam.gserviceaccount.com
gcloud services enable dataflow.googleapis.com compute.googleapis.com storage.googleapis.com
Paso 1: bucket para el resultado y los ficheros temporales
gcloud storage buckets create $BUCKET --location=$REGION --uniform-bucket-level-access
Paso 2: cuenta de servicio de los workers
Por defecto, los workers usan la cuenta de servicio predeterminada de Compute Engine, que en muchos proyectos tiene el rol de Editor: demasiados permisos. Crea una cuenta dedicada con lo justo: roles/dataflow.worker en el proyecto y permiso sobre tu bucket (el bucket de entrada y el de plantillas son públicos).
gcloud iam service-accounts create $SA --display-name="Workers de Dataflow lab 13"
gcloud projects add-iam-policy-binding $PROJECT_ID \
--member=serviceAccount:$SA_EMAIL --role=roles/dataflow.worker
gcloud storage buckets add-iam-policy-binding $BUCKET \
--member=serviceAccount:$SA_EMAIL --role=roles/storage.objectAdmin
Como propietario del proyecto ya puedes asignar esta cuenta al job (iam.serviceAccounts.actAs).
Paso 3: lanzar la plantilla
Por consola: Dataflow → Trabajos → Crear trabajo a partir de plantilla, nombre wordcount-lab13, región europe-southwest1, plantilla Word Count, archivo de entrada gs://dataflow-samples/shakespeare/kinglear.txt, ubicación de salida gs://TU_BUCKET/results/output, ubicación temporal gs://TU_BUCKET/temp y, en parámetros opcionales, la cuenta de servicio.
Con gcloud:
gcloud dataflow jobs run wordcount-lab13 \
--gcs-location=gs://dataflow-templates/latest/Word_Count \
--region=$REGION \
--service-account-email=$SA_EMAIL \
--staging-location=$BUCKET/temp \
--max-workers=2 \
--parameters=inputFile=gs://dataflow-samples/shakespeare/kinglear.txt,output=$BUCKET/results/output
--max-workers=2 pone un techo al autoescalado: es una buena práctica de control de costes en entornos de pruebas.
Paso 4: observar el job
Mientras se ejecuta (tarda unos minutos, la mayor parte en arrancar los workers):
gcloud dataflow jobs list --region=$REGION --status=active
export JOB_ID=$(gcloud dataflow jobs list --region=$REGION --filter="name=wordcount-lab13" --format="value(id)" --limit=1)
gcloud dataflow jobs describe $JOB_ID --region=$REGION --format="yaml(currentState,type,createTime)"
# Los workers son VMs normales de Compute Engine creadas y gestionadas por Dataflow
gcloud compute instances list --filter="name~wordcount"
En la consola (Dataflow → Trabajos → wordcount-lab13) fíjate en:
- Gráfico del trabajo: cada caja es una transformación de Beam (lectura,
ParDoque separa palabras, recuento conGroupByKey/Combine, escritura). Al hacer clic ves el tiempo y los elementos procesados en cada una. - Métricas del trabajo: número de workers a lo largo del tiempo (autoescalado), rendimiento y uso de CPU.
- Información del trabajo: tipo (Batch), región, cuenta de servicio y, si aplica, si usa Dataflow Prime.
- Registros: los logs del job y de los workers en Cloud Logging, lo primero que miras si un job falla (por ejemplo, por permisos de la cuenta de servicio o por la red).
- Coste: la pestaña de coste (cuando está disponible) muestra una estimación del gasto del job.
Cuando el estado sea JOB_STATE_DONE, los workers se habrán eliminado solos: comprueba que gcloud compute instances list --filter="name~wordcount" ya no devuelve nada.
Paso 5: ver el resultado
gcloud storage ls $BUCKET/results/
gcloud storage cat "$BUCKET/results/output*" | sort -t: -k2 -n -r | head -20
Verás las palabras más frecuentes de El rey Lear con su recuento (el formato es palabra: número).
Comprueba que funciona
gcloud dataflow jobs describemuestracurrentState: JOB_STATE_DONEytype: JOB_TYPE_BATCH.- En
results/hay uno o varios ficherosoutput-XXXXX-of-YYYYY(uno por fragmento que escribió cada worker). - No queda ninguna VM de workers en ejecución.
Limpieza
# Por si el job siguiera activo (no debería en un job por lotes terminado)
gcloud dataflow jobs cancel $JOB_ID --region=$REGION 2>/dev/null
gcloud storage rm --recursive $BUCKET
gcloud projects remove-iam-policy-binding $PROJECT_ID \
--member=serviceAccount:$SA_EMAIL --role=roles/dataflow.worker
gcloud iam service-accounts delete $SA_EMAIL --quiet
Comprueba en Dataflow → Trabajos que no hay ningún job en estado Running. Un job en streaming nunca termina solo y factura sin parar: si en algún momento lanzas uno (por ejemplo, la plantilla Pub/Sub a BigQuery), páralo siempre con Drain o Cancel.
Preguntas para pensar como arquitecto
Si lanzaras la plantilla «Pub/Sub Subscription to BigQuery» en lugar de WordCount, ¿qué cambiaría en el comportamiento y en el coste?
Sería un job en streaming: no termina nunca, mantiene al menos un worker encendido de forma permanente y factura 24 h al día. Para pararlo usarías Drain (termina de procesar lo que está en vuelo y no pierde datos) en lugar de Cancel. Si los mensajes no necesitan transformación, una suscripción BigQuery de Pub/Sub haría lo mismo sin coste de workers.
¿Por qué creaste una cuenta de servicio dedicada en lugar de usar la predeterminada de Compute Engine?
Por mínimo privilegio: la cuenta predeterminada suele tener el rol de Editor en todo el proyecto, así que cualquier código que corra en los workers podría modificar casi cualquier recurso. La dedicada solo puede actuar como worker de Dataflow y escribir en un bucket concreto. Es una recomendación explícita de seguridad que el examen valora.
Un job por lotes nocturno de Dataflow no es urgente y el equipo quiere abaratarlo. ¿Qué opción de Dataflow propones?
FlexRS (Flexible Resource Scheduling): Dataflow retrasa el arranque del job (dentro de una ventana) y usa una mezcla de VMs preemptibles y estándar a un precio menor. Es adecuado para lotes que toleran esperar. Además, limita max-workers y revisa el tipo de máquina de los workers.
El equipo tiene 200 jobs Spark on-premises y pregunta si debe reescribirlos en Beam para usar Dataflow. ¿Qué le aconsejas?
No reescribirlos de entrada: migrarlos con cambios mínimos a Managed Service for Apache Spark (antes Dataproc), con clústeres efímeros o en modo serverless y los datos en Cloud Storage. Dataflow tiene sentido para canalizaciones nuevas, especialmente en streaming o cuando se quiere un servicio sin clúster con el mismo código para lotes y streaming. Reescribir 200 jobs tiene un coste y un riesgo que el negocio debe justificar.