SHA256
Feature: CSV|Retraso,FFT,SMA|WebSocket
This commit is contained in:
1 parent
f5a98d9e46
commit
14e3065f60
6 files changed
+6255
-7
No files matched your search
@@ -5,6 +5,8 @@ import queue
|
|||||||
import threading
|
import threading
|
||||||
|
|
||||||
from src.cliente_ws import websocket_cliente
|
from src.cliente_ws import websocket_cliente
|
||||||
|
from src.gurdar import guardar_datos
|
||||||
|
from src.procesor_data import data_procesor
|
||||||
|
|
||||||
# Datos globals
|
# Datos globals
|
||||||
ip_server = "10.42.0.100:81"
|
ip_server = "10.42.0.100:81"
|
||||||
@@ -12,19 +14,31 @@ ip_server = "10.42.0.100:81"
|
|||||||
# Cola de datos compartidos
|
# Cola de datos compartidos
|
||||||
raw_data = queue.Queue()
|
raw_data = queue.Queue()
|
||||||
data_clean = queue.Queue()
|
data_clean = queue.Queue()
|
||||||
|
data_grafica = queue.Queue()
|
||||||
|
i_paquetes = queue.Queue()
|
||||||
|
buffer_size = queue.Queue()
|
||||||
|
|
||||||
|
|
||||||
def debug_print():
|
# def debug_print():
|
||||||
while True:
|
# while True:
|
||||||
raw = raw_data.get()
|
# data = data_clean.get()
|
||||||
print(raw)
|
# # raw = data_grafica.get()
|
||||||
|
# print(data)
|
||||||
|
# # print(raw)
|
||||||
|
# # print(i)
|
||||||
|
|
||||||
|
|
||||||
hilo_1 = websocket_cliente(raw_data, ip_server)
|
hilo_1 = websocket_cliente(raw_data, i_paquetes, ip_server)
|
||||||
hilo_2 = threading.Thread(target=debug_print)
|
hilo_2 = data_procesor(raw_data, data_clean, data_grafica, buffer_size)
|
||||||
|
# hilo_3 = threading.Thread(target=debug_print)
|
||||||
|
hilo_4 = guardar_datos(data_clean)
|
||||||
|
|
||||||
hilo_1.start()
|
hilo_1.start()
|
||||||
hilo_2.start()
|
hilo_2.start()
|
||||||
|
# hilo_3.start()
|
||||||
|
hilo_4.start()
|
||||||
|
|
||||||
hilo_1.join()
|
hilo_1.join()
|
||||||
hilo_2.join()
|
hilo_2.join()
|
||||||
|
# hilo_3.join()
|
||||||
|
hilo_4.join()
|
||||||
+7
-1
@@ -5,18 +5,24 @@ from websocket import create_connection
|
|||||||
|
|
||||||
|
|
||||||
class websocket_cliente(threading.Thread):
|
class websocket_cliente(threading.Thread):
|
||||||
def __init__(self, cola_datos, ip_server):
|
def __init__(self, cola_datos, cola_contador, ip_server):
|
||||||
super().__init__()
|
super().__init__()
|
||||||
self.cola_datos = cola_datos
|
self.cola_datos = cola_datos
|
||||||
self.ip_server = ip_server
|
self.ip_server = ip_server
|
||||||
|
self.cola_contador = cola_contador
|
||||||
self.ws = create_connection(f"ws://{self.ip_server}")
|
self.ws = create_connection(f"ws://{self.ip_server}")
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
|
# variable local en zero - contador
|
||||||
|
self.contador = 0
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
data = self.ws.recv()
|
data = self.ws.recv()
|
||||||
|
|
||||||
if isinstance(data, bytes):
|
if isinstance(data, bytes):
|
||||||
samples = struct.unpack(f"<{len(data) // 2}H", data)
|
samples = struct.unpack(f"<{len(data) // 2}H", data)
|
||||||
self.cola_datos.put(samples)
|
self.cola_datos.put(samples)
|
||||||
|
self.contador += 1
|
||||||
|
self.cola_contador.put(self.contador)
|
||||||
else:
|
else:
|
||||||
print("Datos: TXT. No se procesan.")
|
print("Datos: TXT. No se procesan.")
|
||||||
@@ -0,0 +1,68 @@
|
|||||||
|
import queue
|
||||||
|
import threading
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pyqtgraph as pg
|
||||||
|
from PyQt5 import QtWidgets
|
||||||
|
from pyqtgraph.Qt import QtCore
|
||||||
|
|
||||||
|
|
||||||
|
def graficar(cola_clean, cola_filtrada):
|
||||||
|
app = QtWidgets.QApplication([])
|
||||||
|
win = pg.GraphicsLayoutWidget(show=True, title="Onda senoidal")
|
||||||
|
|
||||||
|
# Un solo plot con ambas curvas
|
||||||
|
p1 = win.addPlot(title="Señal cruda y filtrada")
|
||||||
|
p1.addLegend() # leyenda automática
|
||||||
|
curve1 = p1.plot(pen="r", name="Cruda")
|
||||||
|
curve2 = p1.plot(pen="b", name="Filtrada")
|
||||||
|
|
||||||
|
def update():
|
||||||
|
# Actualizar señal cruda
|
||||||
|
if not cola_clean.empty():
|
||||||
|
datos = cola_clean.get()
|
||||||
|
muestras = np.arange(len(datos))
|
||||||
|
curve1.setData(muestras, datos)
|
||||||
|
|
||||||
|
# Actualizar señal filtrada
|
||||||
|
if not cola_filtrada.empty():
|
||||||
|
datos_f = cola_filtrada.get()
|
||||||
|
muestras_f = np.arange(len(datos_f))
|
||||||
|
curve2.setData(muestras_f, datos_f)
|
||||||
|
|
||||||
|
timer = QtCore.QTimer()
|
||||||
|
timer.timeout.connect(update)
|
||||||
|
timer.start(100) # refresco cada 100 ms
|
||||||
|
|
||||||
|
app.exec_()
|
||||||
|
|
||||||
|
|
||||||
|
class websocket_cliente(threading.Thread):
|
||||||
|
def __init__(self, cola_clean, cola_filtada, buffer_size):
|
||||||
|
super().__init__()
|
||||||
|
self.cola_clean = cola_clean
|
||||||
|
self.cola_filtrado = cola_filtada
|
||||||
|
self.buffer_size = buffer_size
|
||||||
|
self.buffer_filtrado = [0] * self.buffer_size
|
||||||
|
self.buffer_limpio = [0] * self.buffer_size
|
||||||
|
self.idx_despla = self.buffer_size / 100
|
||||||
|
|
||||||
|
def run(self):
|
||||||
|
while True:
|
||||||
|
self.buffer_limpio.append(self.cola_clean.get())
|
||||||
|
try:
|
||||||
|
self.buffer_filtrado = self.cola_filtrado.get_nowait()
|
||||||
|
except queue.Empty:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
# tiempo real
|
||||||
|
if __name__ == "__main__":
|
||||||
|
cola_clean = queue.Queue()
|
||||||
|
cola_filtrada = queue.Queue()
|
||||||
|
|
||||||
|
cliente = websocket_cliente(cola_clean, cola_filtrada, buffer_size=1024)
|
||||||
|
cliente.daemon = True
|
||||||
|
cliente.start()
|
||||||
|
|
||||||
|
graficar(cola_clean, cola_filtrada)
|
||||||
@@ -0,0 +1,26 @@
|
|||||||
|
import csv
|
||||||
|
import threading
|
||||||
|
|
||||||
|
|
||||||
|
class guardar_datos(threading.Thread):
|
||||||
|
def __init__(self, cola_datos):
|
||||||
|
super().__init__()
|
||||||
|
self.cola_datos = cola_datos
|
||||||
|
self.daemon = (
|
||||||
|
True # Permite que el hilo cierre si el programa principal termina
|
||||||
|
)
|
||||||
|
|
||||||
|
def run(self):
|
||||||
|
# Abrimos el archivo en modo append ('a')
|
||||||
|
with open("datos.csv", "a", newline="") as archivo_csv:
|
||||||
|
escritor_csv = csv.writer(archivo_csv)
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
# Espera un elemento de la cola
|
||||||
|
fila = self.cola_datos.get()
|
||||||
|
escritor_csv.writerow(fila)
|
||||||
|
# Fuerza la escritura
|
||||||
|
archivo_csv.flush()
|
||||||
|
except Exception as e:
|
||||||
|
print(f"Error al guardar datos: {e}")
|
||||||
|
break
|
||||||
@@ -0,0 +1,73 @@
|
|||||||
|
import threading
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
|
||||||
|
def bit_voltio(raw_data):
|
||||||
|
cons_adc = 0.00080586
|
||||||
|
volts = [cons_adc * lectura for lectura in raw_data]
|
||||||
|
return volts
|
||||||
|
|
||||||
|
|
||||||
|
def procesar(samples, T):
|
||||||
|
SIZE_BUFFER = len(samples)
|
||||||
|
fft = np.fft.fft(samples)
|
||||||
|
frecuencia = np.fft.fftfreq(SIZE_BUFFER, T)
|
||||||
|
|
||||||
|
indices_fft = np.where(frecuencia > 0)
|
||||||
|
f_positivo = frecuencia[indices_fft]
|
||||||
|
magnitudes = np.abs(fft[indices_fft])
|
||||||
|
f_dominante = f_positivo[np.argmax(magnitudes)]
|
||||||
|
array_data = pd.Series(samples)
|
||||||
|
window_size = max(1, int(f_dominante))
|
||||||
|
media_movil = array_data.rolling(window=window_size).mean()
|
||||||
|
return media_movil
|
||||||
|
|
||||||
|
|
||||||
|
def calcular_desfase(original, filtrada, periodo):
|
||||||
|
# Centrar señales en cero
|
||||||
|
orig_centrada = original - np.mean(original)
|
||||||
|
filt_centrada = filtrada - np.mean(filtrada)
|
||||||
|
# Encontrar la coincidencia de fase: Correlación cruzada
|
||||||
|
correlacion = np.correlate(orig_centrada, filt_centrada, mode="full")
|
||||||
|
# Encontrar la distancia respecto al centro
|
||||||
|
centro = len(orig_centrada) - 1
|
||||||
|
desfase_muestras = np.argmax(correlacion) - centro
|
||||||
|
# Segundos
|
||||||
|
desfase_tiempo = desfase_muestras * periodo
|
||||||
|
return desfase_tiempo
|
||||||
|
|
||||||
|
|
||||||
|
class data_procesor(threading.Thread):
|
||||||
|
bandera_analisis = False
|
||||||
|
buffer_fft = []
|
||||||
|
|
||||||
|
def __init__(self, cola_raw, cola_clean, cola_graficar, cola_tamaño):
|
||||||
|
super().__init__()
|
||||||
|
self.cola_raw = cola_raw
|
||||||
|
self.cola_clean = cola_clean
|
||||||
|
self.cola_graficar = cola_graficar
|
||||||
|
self.cola_tamaño = cola_tamaño
|
||||||
|
self.size_array = 0
|
||||||
|
|
||||||
|
def run(self):
|
||||||
|
self.periodo = 0
|
||||||
|
while True:
|
||||||
|
raw = self.cola_raw.get()
|
||||||
|
data_procesor.buffer_fft.extend(raw)
|
||||||
|
if not data_procesor.bandera_analisis:
|
||||||
|
self.size_array = len(raw)
|
||||||
|
self.periodo = 1 / (self.size_array * 100)
|
||||||
|
self.total_size = 100 * self.size_array
|
||||||
|
self.cola_tamaño.put(self.total_size)
|
||||||
|
if len(data_procesor.buffer_fft) >= self.total_size:
|
||||||
|
senal_original = np.array(data_procesor.buffer_fft)
|
||||||
|
data_clean = procesar(data_procesor.buffer_fft, self.periodo)
|
||||||
|
self.cola_graficar.put(data_clean)
|
||||||
|
tiempo = calcular_desfase(senal_original, data_clean, self.periodo)
|
||||||
|
print(f"Desface: {tiempo:.6f}s")
|
||||||
|
data_procesor.buffer_fft.clear()
|
||||||
|
|
||||||
|
# DEBUG print(self.size_array)
|
||||||
|
self.cola_clean.put(bit_voltio(raw))
|
||||||
Reference in new issue
Block a user