Lab práctico · Semana 6: Procesamiento de datos y analítica

Ingesta de eventos con Pub/Sub y suscripción BigQuery

⏱ 60-90 minDificultad: mediaApartados: 2.21.3

Qué vas a construir

Una canalización de ingesta sin código ni servidores: simulas sensores que publican lecturas en un tema de Pub/Sub, y una suscripción BigQuery las escribe directamente en una tabla particionada. Después creas una tabla curada particionada y agrupada, consultas los datos y aprendes a ver y limitar el coste de las consultas antes de ejecutarlas.

flowchart LR
  S["Sensores simulados (Cloud Shell)"] -->|"gcloud pubsub topics publish"| T["Tema lecturas"]
  T --> SUB["Suscripción BigQuery lecturas-bq"]
  SUB --> RAW["iot.lecturas_raw (particionada por publish_time)"]
  RAW -->|"CREATE TABLE AS SELECT"| CUR["iot.lecturas (PARTITION BY fecha, CLUSTER BY sensor_id)"]
  CUR --> Q["Consultas, dry run y máximo de bytes facturados"]

Antes de empezar

  • Proyecto dedicado con facturación y Cloud Shell abierto.
  • Conviene haber leído las secciones de Pub/Sub y de BigQuery del módulo 6.
export PROJECT_ID=$(gcloud config get-value project)
export PROJECT_NUMBER=$(gcloud projects describe $PROJECT_ID --format="value(projectNumber)")
export REGION=europe-southwest1
gcloud services enable pubsub.googleapis.com bigquery.googleapis.com

Paso 1: dataset y tabla de aterrizaje en BigQuery

La suscripción BigQuery necesita que la tabla exista. Usarás la opción escribir metadatos (write metadata), que exige estas columnas: subscription_name, message_id, publish_time, data y attributes. El mensaje (JSON) va a la columna data de tipo JSON. Particionas por publish_time para que las consultas por fecha lean solo las particiones necesarias.

Crea el dataset en Madrid y el fichero de esquema:

bq --location=$REGION mk --dataset ${PROJECT_ID}:iot

cat > esquema_raw.json <<'EOF'
[
  {"name": "subscription_name", "type": "STRING"},
  {"name": "message_id", "type": "STRING"},
  {"name": "publish_time", "type": "TIMESTAMP"},
  {"name": "data", "type": "JSON"},
  {"name": "attributes", "type": "JSON"}
]
EOF

bq mk --table \
  --time_partitioning_field=publish_time \
  --time_partitioning_type=DAY \
  ${PROJECT_ID}:iot.lecturas_raw ./esquema_raw.json

bq show --format=prettyjson ${PROJECT_ID}:iot.lecturas_raw | grep -A3 timePartitioning

Por consola: BigQuery → Studio → tu proyecto → Crear conjunto de datos (ubicación europe-southwest1) y después Crear tabla con el esquema anterior y la partición por el campo publish_time.

Paso 2: tema y permisos de la cuenta de servicio de Pub/Sub

Quien escribe en BigQuery no eres tú, sino el agente de servicio de Pub/Sub (service-NÚMERO@gcp-sa-pubsub.iam.gserviceaccount.com). Necesita roles/bigquery.dataEditor. Para respetar el mínimo privilegio, se lo das solo sobre el dataset iot, no sobre todo el proyecto.

gcloud pubsub topics create lecturas

# Asegura que existe el agente de servicio de Pub/Sub
gcloud beta services identity create --service=pubsub.googleapis.com --project=$PROJECT_ID

export PUBSUB_SA=service-${PROJECT_NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.com

Para conceder el rol sobre el dataset, lo más sencillo es la consola: BigQuery → dataset iot → Compartir → Permisos → Añadir principal, pega el valor de $PUBSUB_SA y elige BigQuery Data Editor. Si prefieres la línea de comandos, con una sentencia DCL de BigQuery:

bq query --use_legacy_sql=false \
  "GRANT \`roles/bigquery.dataEditor\` ON SCHEMA \`${PROJECT_ID}.iot\` TO 'serviceAccount:${PUBSUB_SA}'"

Paso 3: suscripción BigQuery

gcloud pubsub subscriptions create lecturas-bq \
  --topic=lecturas \
  --bigquery-table=${PROJECT_ID}:iot.lecturas_raw \
  --write-metadata

gcloud pubsub subscriptions describe lecturas-bq --format="yaml(bigqueryConfig,state)"

state: ACTIVE indica que Pub/Sub puede escribir en la tabla. Si aparece un error de permisos, espera un par de minutos (propagación de IAM) y revisa el paso 2.

Paso 4: publicar lecturas de sensores

Publica 60 mensajes JSON de 10 sensores con temperaturas aleatorias y un atributo con la planta (tarda uno o dos minutos):

for i in $(seq 1 60); do
  SENSOR="sensor-$(( RANDOM % 10 + 1 ))"
  TEMP="$(( RANDOM % 15 + 15 )).$(( RANDOM % 10 ))"
  AHORA=$(date -u +%Y-%m-%dT%H:%M:%SZ)
  gcloud pubsub topics publish lecturas \
    --message="{\"sensor_id\":\"${SENSOR}\",\"temperatura\":${TEMP},\"ts\":\"${AHORA}\"}" \
    --attribute=planta=madrid > /dev/null
done
echo "Publicados"

Paso 5: consultar los datos

Los mensajes aparecen en BigQuery en segundos:

bq query --use_legacy_sql=false '
SELECT
  JSON_VALUE(data, "$.sensor_id") AS sensor,
  COUNT(*) AS lecturas,
  ROUND(AVG(CAST(JSON_VALUE(data, "$.temperatura") AS FLOAT64)), 1) AS media
FROM iot.lecturas_raw
WHERE DATE(publish_time) = CURRENT_DATE()
GROUP BY sensor
ORDER BY sensor'

Y comprueba los metadatos que ha escrito Pub/Sub:

bq query --use_legacy_sql=false \
  'SELECT message_id, publish_time, attributes FROM iot.lecturas_raw LIMIT 3'

Paso 6: tabla curada particionada y agrupada

La tabla de aterrizaje guarda el JSON en bruto. Para analítica, crea una tabla con columnas tipadas, particionada por el día de la lectura y agrupada por sensor (patrón ELT: cargas en bruto y transformas con SQL dentro de BigQuery):

bq query --use_legacy_sql=false '
CREATE TABLE iot.lecturas
PARTITION BY DATE(ts)
CLUSTER BY sensor_id
OPTIONS (require_partition_filter = TRUE)
AS
SELECT
  JSON_VALUE(data, "$.sensor_id") AS sensor_id,
  CAST(JSON_VALUE(data, "$.temperatura") AS FLOAT64) AS temperatura,
  TIMESTAMP(JSON_VALUE(data, "$.ts")) AS ts,
  publish_time
FROM iot.lecturas_raw
WHERE DATE(publish_time) = CURRENT_DATE()'

Prueba la opción require_partition_filter: esta consulta falla porque no filtra por la columna de partición, lo que protege frente a escaneos completos accidentales:

bq query --use_legacy_sql=false 'SELECT COUNT(*) FROM iot.lecturas'

Y esta funciona:

bq query --use_legacy_sql=false \
  'SELECT sensor_id, MAX(temperatura) AS maxima FROM iot.lecturas WHERE DATE(ts) = CURRENT_DATE() GROUP BY sensor_id'

Paso 7: ver el coste de las consultas antes de ejecutarlas

Con datos tan pequeños todo se factura al mínimo (10 MB por tabla) y cabe en el TiB gratuito mensual. Para ver el efecto real del diseño, usa el dry run sobre un dataset público grande: BigQuery calcula los bytes que escanearía sin ejecutar ni cobrar nada.

# Todas las columnas
bq query --use_legacy_sql=false --dry_run \
  'SELECT * FROM `bigquery-public-data.stackoverflow.posts_questions`'

# Solo una columna
bq query --use_legacy_sql=false --dry_run \
  'SELECT title FROM `bigquery-public-data.stackoverflow.posts_questions`'

Compara los bytes: el almacenamiento columnar hace que leer una columna cueste una fracción de SELECT *. Multiplica por 6,25 USD/TiB (precio bajo demanda de lista en EE. UU.) para estimar el coste.

Ahora pon un límite de bytes facturados: la consulta falla antes de ejecutarse, sin coste:

bq query --use_legacy_sql=false --maximum_bytes_billed=10000000 \
  'SELECT title FROM `bigquery-public-data.stackoverflow.posts_questions`'

Por último, revisa lo que han procesado y facturado tus consultas del lab en la vista de trabajos de la región:

bq query --use_legacy_sql=false '
SELECT creation_time, total_bytes_processed, total_bytes_billed, LEFT(query, 60) AS consulta
FROM `region-europe-southwest1`.INFORMATION_SCHEMA.JOBS_BY_USER
WHERE creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 2 HOUR)
ORDER BY creation_time DESC'

En la consola, el validador (arriba a la derecha del editor) muestra la misma estimación de bytes mientras escribes la consulta, y la pestaña Información del trabajo muestra los bytes facturados.

Comprueba que funciona

  • La suscripción lecturas-bq está ACTIVE y la tabla iot.lecturas_raw tiene unas 60 filas con message_id, publish_time y attributes rellenos.
  • La consulta del paso 5 devuelve la media por sensor.
  • iot.lecturas aparece en la consola como tabla particionada por DAY y agrupada por sensor_id, y rechaza consultas sin filtro de partición.
  • El dry run de SELECT * estima muchos más bytes que el de una sola columna, y la consulta con --maximum_bytes_billed falla sin coste.

Limpieza

gcloud pubsub subscriptions delete lecturas-bq
gcloud pubsub topics delete lecturas
bq rm -r -f -d ${PROJECT_ID}:iot
rm -f esquema_raw.json

Al borrar el dataset desaparece también el permiso que diste al agente de servicio de Pub/Sub sobre él.

Preguntas para pensar como arquitecto

¿Cuándo sustituirías la suscripción BigQuery por Dataflow?

Cuando haya que transformar los eventos antes de guardarlos: agregar por ventanas de tiempo, enriquecer con otras fuentes, deduplicar, tratar datos que llegan tarde o escribir en varios destinos (Bigtable y BigQuery a la vez). Si los mensajes ya llegan listos para guardarse, la suscripción BigQuery es más simple, más barata y no requiere operar ninguna canalización; las transformaciones pueden hacerse después con SQL (ELT).

Algunos mensajes tienen un formato incorrecto y no se pueden escribir en la tabla. ¿Qué pasa y cómo lo diseñarías?

Un mensaje que no se puede escribir en BigQuery no se confirma, así que Pub/Sub lo reintenta. Configura un dead letter topic en la suscripción (con un máximo de intentos entre 5 y 100) para apartar esos mensajes en otro tema, conservarlos (por ejemplo, con una suscripción Cloud Storage) y analizarlos sin bloquear el resto. El agente de servicio de Pub/Sub necesita permiso para publicar en el tema de dead letter y para suscribirse a la suscripción original.

El equipo de datos lanza consultas SELECT * sobre una tabla de 50 TB varias veces al día y la factura se ha disparado. Propón tres medidas.
  1. Particionar la tabla por fecha y agruparla por las columnas de filtro, con require_partition_filter. 2) Enseñar a seleccionar solo las columnas necesarias y fijar maximum_bytes_billed o cuotas personalizadas por usuario o proyecto. 3) Si el gasto es alto y constante, pasar a precios por capacidad (edición con reservas y compromiso) para tener un coste predecible. Además, vistas materializadas para agregados repetidos.
¿Por qué el dataset se crea en europe-southwest1 y no en la multirregión EU?

Porque la ubicación del dataset determina dónde se almacenan y procesan los datos. Si el requisito es residencia de datos en España o minimizar la latencia con otros recursos de Madrid, una región concreta es lo correcto. La multirregión EU aporta más redundancia geográfica pero reparte los datos por varios países de la UE. Recuerda que no puedes hacer joins entre datasets de ubicaciones distintas en una misma consulta.


Volver al módulo