Recolector de datos de una pesadora / control de calidad ULMA (gateway CONTEC CPS‑MCS341‑DS configurado por ANRITSU) mediante OPC UA. Guarda el dato crudo en ficheros, consolida una línea por producción en PostgreSQL y expone KPIs de ritmo (kg/hora) para un panel.
Proyecto 520746.
Pesadora ULMA (OPC UA)
│ ns=2;s=CPS-MCS341-DS.STAG.*
▼
┌───────────────────────────────────────────────┐
│ RECOGER (Python · asyncio · main.py) │
│ - Suscripción OPC UA (callback no bloqueante) │
│ - Encola cada evento en 3 colas │
└───────────────────────────────────────────────┘
│ │ │
▼ ▼ ▼
EVENT_QUEUE STATS_QUEUE DB_QUEUE
(asyncio) (queue, hilo) (queue, hilo)
│ │ │
▼ ▼ ▼
GUARDAR: MM-all.jsonl MM-stats.json PostgreSQL
(crudo, todo) (crudo stats) (pesadora_lineas, consolidado)
│
▼
MOSTRAR: API PHP → ritmo kg/hora (panel)
Idea central: recoger ≠ guardar ≠ mostrar. Cada capa cambia por motivos distintos y no se estorban.
- La pesadora publica un cambio de tag (
STAGxx) por OPC UA. SusctiptionHandler.datachange_notification(ensuscription_hadler_file.py) construye uneventy solo encola (trabajo mínimo, no bloquea).- Consumidores independientes escriben:
EVENT_QUEUE→jsonl_writer→AAAA/MM-all.jsonl(todos los tags).STATS_QUEUE→stats_writer→AAAA/MM-stats.json(tags de estadística).index_app(app_recalculate_file.py) acumula el estado y, cuando avanza el peso y los estadísticos son coherentes, encola enDB_QUEUE.DB_QUEUE→database_file→INSERTen PostgreSQL (idempotente).
- Un heartbeat en hilo aparte escribe periódicamente si el proceso sigue vivo y conectado.
main.py Arranque: suscripción OPC UA + lanzamiento de workers
config.py Configuración desde .env (falla rápido si falta algo)
requirements.txt
assets/
suscription_hadler_file.py Callback OPC UA -> encola eventos
event_queue_file.py EVENT_QUEUE (asyncio) · STATS_QUEUE, DB_QUEUE (queue)
jsonl_writer.py Consumidor async -> MM-all.jsonl
stats_writer_file.py Hilo -> MM-stats.json
app_recalculate_file.py Consolidación por producción -> DB_QUEUE
database_file.py Hilo -> INSERT PostgreSQL (conexión única + reconexión)
hertbeat_writer_file.py Hilo -> heartbeat (proceso vivo / conexión)
conection_state_file.py Estado de conexión (thread-safe)
supervised_file.py Relanza una corrutina si muere
utils_file.py Fechas y parseo de valores
Todos son STRING, solo lectura, patrón ns=2;s=CPS-MCS341-DS.STAG.{tag}.
Estado / artículo en producción
STAG00operación (0=funciona, 1=parado, 2=ajuste, 3=check) ·STAG01estado técnicoSTAG02producto seleccionado ·STAG21/22producto y nombre de las estadísticas
Pesada individual (hasta 70/min; puede colapsar por el filtro de cambio de valor)
STAG10nº secuencial ·STAG11producto ·STAG12clasificación (zona/OK-NG)STAG13rechazo (metal/externo/normal) ·STAG14peso medido (g)
Estadísticas acumuladas (fiables; se reinician por lote)
STAG37muestras buenas ·STAG38peso total (kg) ·STAG39peso medio (g)STAG53conteo total de bolsas ·STAG28–36zonas y rechazosSTAG23/24lote/batch ·STAG25/26inicio/fin ·STAG27método (0=todas / 1=solo PASS)
Recuento oficial = estadísticas (
STAG53, zonas), no el stream por pesada.
python3 -m venv .venv
source .venv/bin/activate
pip install -r requirements.txt # asyncua, psycopg, python-dotenvOPCUA_URL=opc.tcp://xxx:4840
OPCUA_NODE_ID_PREFIX=ns=2;s=CPS-MCS341-DS.STAG.
JSONL_BASE_DIRECTORY=/var/lib/suz-opcua
HEARTBEAT_FILE_NAME=head-bit.json
POSTGRES_HOST=127.0.0.1
POSTGRES_PORT=5432
POSTGRES_DB=suz
POSTGRES_USER=suz
POSTGRES_PASSWORD=********CREATE TABLE public.pesadora_lineas (
id bigserial PRIMARY KEY,
art_erp text,
art_name text,
lote text,
batch text,
inicio_of timestamp,
fin_of timestamp,
bolsas_buenas integer,
kg numeric(12,3),
peso_medio numeric(10,2),
bolsas_total integer,
"date" timestamp
);
-- Idempotencia: evita duplicados aunque se reinicie o arranque dos veces
ALTER TABLE public.pesadora_lineas
ADD CONSTRAINT pesadora_lineas_inicio_of_kg_unique UNIQUE (inicio_of, kg);Cada línea se guarda cuando avanza el peso acumulado y los estadísticos son
coherentes (kg·1000 / bolsas_buenas ≈ peso_medio, tolerancia 0,1 g). Así solo se
persisten snapshots válidos.
/etc/systemd/system/suz-opcua.service:
[Unit]
Description=suz OPC UA collector
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=suz
WorkingDirectory=/opt/suz-opcua
ExecStart=/opt/suz-opcua/.venv/bin/python /opt/suz-opcua/main.py
Environment=PYTHONUNBUFFERED=1
StateDirectory=suz-opcua
Restart=always
RestartSec=5
[Install]
WantedBy=multi-user.targetsudo systemctl daemon-reload
sudo systemctl enable --now suz-opcua
sudo systemctl restart suz-opcua
journalctl -u suz-opcua -fRitmo kg/hora en ventanas de 5 / 25 / 60 min: suma el incremento de kg del tramo
continuo del artículo actual y lo lleva a hora.
Notas importantes:
- Poner
date_default_timezone_set('Europe/Madrid')(o anclar la ventana aMAX(fin_of)) porquefin_ofes hora de la máquina y el servidor puede ir en UTC. - Al reiniciarse el
kgpor lote, sikg_actual < kg_anterior= nuevo lote: no cruzar esa frontera al calcular.
- No bloquear la recogida: el callback OPC solo encola; la I/O va en workers.
- Colas + workers (productor/consumidor):
asyncio.Queuepara consumidores async,queue.Queuepara consumidores en hilo (¡no mezclar!). - Todo se reconecta y se reinicia solo; nada muere en silencio.
- Idempotencia: clave única +
ON CONFLICT DO NOTHING. - Validar antes de persistir (coherencia de los estadísticos).
- Guardar el crudo (JSONL) y derivar lo procesado (PostgreSQL).
- Config por
.env; una sola zona horaria de referencia. - Observabilidad: heartbeat + logs en journald.