Real-Time Big Data & LLM Agents para la Industria

Apache Kafka · ClickHouse · PySpark Streaming · Vector DBs · LangChain/LlamaIndex · Multi-Agent RAG

Real-Time Big Data & LLM Agents Cover
Level 3 · Master Data & AI Architect

Capítulo 1 · Ingesta & Streaming IIoT con Apache Kafka & Sparkplug B

En plantas industriales modernas y redes IoT de alta densidad, se generan millones de eventos por segundo procedentes de sensores de vibración, temperatura, variadores de frecuencia y PLCs. Para ingerir y desacoplar este volumen masivo de datos sin pérdida de mensajes, se utiliza **Apache Kafka** configurado con el protocolo **MQTT Sparkplug B**.

¿Qué es una Arquitectura de Streaming Kappa?
La arquitectura Kappa simplifica el procesamiento reemplazando la capa por lotes (Batch) tradicional por un único log distribuido e inmutable (Kafka), donde los datos históricos y en tiempo real son procesados por el mismo pipeline de streaming.
┌─────────────────────────────────────────────────────────────────────────────┐ │ ARQUITECTURA DE INGESTA IIOT CON KAFKA │ │ │ │ [ Sensores IIoT / PLCs ] ──► [ MQTT Broker ] │ │ │ (MQTT Sparkplug B Engine) │ │ ▼ │ │ ┌──────────────────────┐ │ │ │ Apache Kafka Cluster │ │ │ │ Topic: iiot.telemetry│ │ │ └──────────┬───────────┘ │ │ │ │ │ ┌────────────────────────┴────────────────────────┐ │ │ ▼ ▼ │ │ [ PySpark Stream Engine ] [ ClickHouse Engine ] │ │ (Window Aggregations & Anomaly) (OLAP Real-Time Analytics) │ └─────────────────────────────────────────────────────────────────────────────┘

Configuración de Productor Kafka en Python (librería `confluent-kafka`)

python (Kafka Producer IIoT Telemetry)
from confluent_kafka import Producer
import json, time, random

conf = {
    'bootstrap.servers': 'kafka-cluster.internal:9092',
    'client.id': 'iiot-sensor-gateway-01',
    'acks': '1', # Latencia ultrabaja asegurando recepción por el broker líder
    'compression.type': 'lz4'
}

producer = Producer(conf)

def send_telemetry(sensor_id, temp, vib):
    payload = {
        "timestamp": int(time.time() * 1000),
        "sensor_id": sensor_id,
        "temperature_c": temp,
        "vibration_hz": vib,
        "status": "NORMAL" if temp < 85.0 else "WARNING"
    }
    producer.produce(
        'iiot.telemetry',
        key=sensor_id.encode('utf-8'),
        value=json.dumps(payload).encode('utf-8')
    )
    producer.poll(0)

# Simulación de transmisión en tiempo real
send_telemetry("TURBINE-04", 78.4, 120.5)
producer.flush()
Cuestionario — Capítulo 1
1. En una arquitectura de streaming masiva para IIoT, ¿cuál es la función de particionar (Partitioning) los Topics en Apache Kafka?
Las particiones son la unidad básica de escalabilidad en Kafka. Al particionar por `sensor_id`, garantizamos el orden estricto de eventos por dispositivo mientras escalamos el consumo horizontalmente en paralelo.

Capítulo 2 · Almacenamiento OLAP de Ultra-Baja Latencia (ClickHouse & TimescaleDB)

Las bases de datos relacionales tradicionales como PostgreSQL o MySQL colapsan bajo cargas de inserción de 50,000 registros/segundo. **ClickHouse** es un motor columnar OLAP diseñado para ingerir y consultar petabytes de telemetría a sub-segundo.

¿Por qué ClickHouse en la Industria 4.0? Su motor de almacenamiento columnar comprime datos de series temporales hasta un 90% (usando códecs como `DoubleDelta` y `Gorilla`), permitiendo agregaciones vectorizadas SIMD instantáneas.

Creación de Tabla Optimizada `ReplacingMergeTree` en ClickHouse

sql (ClickHouse Telemetry DDL)
CREATE TABLE industrial.sensor_telemetry
(
    timestamp DateTime64(3, 'UTC'),
    sensor_id LowCardinality(String),
    temperature_c Float32 CODEC(Gorilla, ZSTD),
    vibration_hz Float32 CODEC(DoubleDelta, ZSTD),
    pressure_bar Float32,
    status Enum8('NORMAL' = 1, 'WARNING' = 2, 'CRITICAL' = 3)
)
ENGINE = ReplacingMergeTree(timestamp)
PRIMARY KEY (sensor_id, timestamp)
ORDER BY (sensor_id, timestamp);

-- Consulta de Agregación en Tiempo Real sobre 100 Millones de Eventos
SELECT 
    toStartOfMinute(timestamp) AS minute,
    sensor_id,
    avg(temperature_c) AS avg_temp,
    max(vibration_hz) AS max_vibration
FROM industrial.sensor_telemetry
WHERE timestamp >= now() - INTERVAL 1 HOUR
GROUP BY minute, sensor_id
ORDER BY minute DESC;
Cuestionario — Capítulo 2
1. ¿Qué ventaja ofrece la estructura de almacenamiento columnar de ClickHouse sobre las bases de datos orientadas a filas?
En bases de datos columnares, los datos de la misma columna se almacenan contiguos en disco. Esto maximiza la eficiencia de caché CPU, compresión de datos y evita lecturas de disco innecesarias.

Capítulo 3 · PySpark Structured Streaming & Anomaly Analytics

Para analizar tendencias y detectar picos térmicos o variaciones armónicas en tiempo real a medida que fluyen los eventos por Kafka, se utilizan ventanas deslizantes (Sliding Windows) en **PySpark Structured Streaming**.

python (PySpark Sliding Window Pipeline)
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window, avg, max

spark = SparkSession.builder \
    .appName("IIoT-RealTime-Anomaly-Detector") \
    .getOrCreate()

# Lectura en directo desde Kafka Topic
kafka_stream = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "iiot.telemetry") \
    .load()

# Definición de Ventana Deslizante de 5 minutos con paso de 1 minuto
windowed_analytics = telemetry_df \
    .groupBy(
        window(col("timestamp"), "5 minutes", "1 minute"),
        col("sensor_id")
    ) \
    .agg(
        avg("temperature_c").alias("avg_temp"),
        max("vibration_hz").alias("peak_vibration")
    )

query = windowed_analytics.writeStream \
    .outputMode("update") \
    .format("console") \
    .start()
Cuestionario — Capítulo 3
1. ¿Cuál es la diferencia principal entre una Tumbling Window y una Sliding Window en PySpark Structured Streaming?
Sliding Windows se solapan continuamente (ej. ventana de 5 min actualizándose cada 10 segundos), permitiendo una detección instantánea de picos antes de que se complete el bloque de tiempo.

Capítulo 4 · RAG Industrial & Vector Databases (Qdrant & Milvus)

Los modelos de lenguaje (LLMs) por sí solos no conocen la configuración específica ni los esquemas P&ID de una planta industrial. Mediante **Retrieval-Augmented Generation (RAG)**, alimentamos al LLM con manuales de mantenimiento, registros de averías y datos de sensores indexados en una **Vector Database (Qdrant)**.

┌─────────────────────────────────────────────────────────────────────────────┐ │ ARQUITECTURA RAG INDUSTRIAL HÍBRIDA │ │ │ │ [ Manuales PDF / P&ID Specs ] ──► [ Text Chunking + BGE Embeddings ] │ │ │ │ │ ▼ │ │ ┌──────────────────────┐ │ │ │ Qdrant Vector Store │ │ │ └──────────┬───────────┘ │ │ │ (Vector Similarity Search) │ │ [ Pregunta Operador ] ───────────────────────┼───────────────────────────┐ │ │ ▼ │ │ │ [ RAG Context Assembly ] │ │ │ │ │ │ │ ▼ ▼ │ │ [ LLM Agent (Llama 3 / Claude) ] ◄────────┘ │ │ │ │ │ ▼ │ │ [ Respuesta Técnica Precisa ] │ └─────────────────────────────────────────────────────────────────────────────┘

Script de Búsqueda Semántica RAG con Qdrant & Python

python (Qdrant Hybrid Search Engine)
from qdrant_client import QdrantClient
from sentence_transformers import SentenceTransformer

client = QdrantClient(host="qdrant.internal", port=6333)
model = SentenceTransformer('BAAI/bge-small-en-v1.5')

def query_industrial_docs(user_query, machine_type="TURBINE"):
    query_vector = model.encode(user_query).tolist()
    
    search_result = client.search(
        collection_name="industrial_manuals",
        query_vector=query_vector,
        query_filter={
            "must": [{"key": "equipment", "match": {"value": machine_type}}]
        },
        limit=3
    )
    return [hit.payload["text_snippet"] for hit in search_result]

docs = query_industrial_docs("Procedimiento de aislamiento por sobrecalentamiento en cojinetes")
print("[+] Contexto RAG Recuperado:", docs[0])
Cuestionario — Capítulo 4
1. En un sistema RAG industrial, ¿por qué es fundamental aplicar filtros de metadatos (ej. `equipment: TURBINE-04`) junto con la búsqueda vectorial?
El filtrado híbrido por metadatos garantiza que el contexto entregado al LLM corresponda exactamente a la marca, modelo y revisión del equipo que presenta la falla.

Capítulo 5 · Multi-Agent Orchestration & ReAct Loops (LangGraph)

Los sistemas complejos de planta requieren múltiples agentes autónomos colaborando: un **Agente de Telemetría** que ejecuta consultas SQL en ClickHouse, un **Agente de Diagnóstico RAG** que analiza manuales técnicos y un **Agente de Protocolos de Seguridad** que valida normativas.

Telemetry Query Agent
Convierte lenguaje natural a consultas ClickHouse/TimescaleDB en tiempo real para verificar el estado de los sensores.
RAG Diagnostic Agent
Consulta vectores en Qdrant para obtener los pasos exactos de mantenimiento correctivo dictados por el fabricante.
Safety Enforcement Agent
Valida que ninguna recomendación técnica viole las normativas ISO 45001 ni sobrepase los límites de tolerancia física de la planta.
Cuestionario — Capítulo 5
1. ¿En qué consiste el patrón de razonamiento ReAct (Reasoning + Acting) en Agentes LLM?
ReAct permite a los agentes razonar dinámicamente, ejecutar herramientas (como consultas SQL a ClickHouse o API calls) e inspeccionar los datos reales devueltos antes de dar una conclusión final.

Capítulo 6 · Simulador Live Telemetry & LLM Copilot Console

Interactúa con la consola en tiempo real de ingestión de telemetría y diagnóstico por Agente LLM Industrial. Genera lecturas de sensores y consulta al copiloto neuronal.

Turbine Temp (°C)
68.4
NORMAL
Vibration (Hz)
110.2
STABLE
Boiler Press (Bar)
12.5
OPTIMAL
[SYS] [KAFKA] Cluster conectado a topic iiot.telemetry. Stream listo.