Real-Time Big Data & LLM Agents para la Industria
Apache Kafka · ClickHouse · PySpark Streaming · Vector DBs · LangChain/LlamaIndex · Multi-Agent RAG
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**.
Configuración de Productor Kafka en Python (librería `confluent-kafka`)
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()
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.
Creación de Tabla Optimizada `ReplacingMergeTree` en ClickHouse
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;
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**.
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()
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)**.
Script de Búsqueda Semántica RAG con Qdrant & Python
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])
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.
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.