¿Manejando dos flujos de datos entrantes y combinándolos en Python?
He estado investigando varias opciones en Python de subprocesamiento, multiprocesamiento asíncrono, etc., como formas de manejar dos flujos entrantes y combinarlos. Hay mucha información sobre, pero los ejemplos a menudo son intrincados y complicados, y más comúnmente consisten en dividir una sola tarea en múltiples subprocesos o procesos para acelerar el resultado final de la tarea.
Tengo un flujo de datos que ingresa a través de un socket (actualmente uso UDP como otra aplicación que se ejecuta localmente en mi PC, pero puedo considerar cambiar a TCP en el futuro si la aplicación debe ejecutarse en una PC separada) y un flujo en serie entrando a través de un adaptador RS232, y necesito combinar las transmisiones. Esta nueva secuencia se retransmite luego en otro socket.
El problema es que llegan a diferentes velocidades (los datos en serie llegan a 125 hz, los datos del socket a 60-120 hz), así que quiero agregar los últimos datos en serie a los datos del socket.
Mi pregunta es esencialmente cuál es la mejor manera de manejar esto, según la experiencia previa de otras personas. Dado que esta es esencialmente una tarea de E / S, se presta más al subprocesamiento (que sé que está limitado a la concurrencia por el GIL), pero debido a la alta tasa de entrada, me pregunto si el procesamiento múltiple es el camino a seguir.
Si usa subprocesos, supongo que la mejor manera de acceder a cada recurso compartido es usar un bloqueo para escribir los datos en serie en un objeto, y en un subproceso separado cada vez que hay nuevos datos de socket y luego adquirir el bloqueo, acceder a los últimos datos en serie en el objeto, procesándolo y luego enviándolo al otro socket. Sin embargo, el hilo principal tiene mucho trabajo entre cada nuevo mensaje de socket entrante.
Con multiprocesamiento, podría usar una tubería para solicitar y recibir los últimos datos en serie del otro proceso, pero eso solo descarga el manejo de datos en serie y aún deja mucho para el proceso principal.
Respuestas
¿Estás seguro de que necesitas subprocesos múltiples aquí? Si no fuera estrictamente necesario, seguro que lo evitaría.
- No he estado programando demasiado últimamente contra puertos serie y sockets, pero hasta donde yo sé, para ambos, los datos están almacenados en búfer por HW / middleware, por lo que desde esa perspectiva no debería haber necesidad de un hilo por flujo entrante.
- con respecto al hilo principal que tiene mucho trabajo por hacer: ¿estás seguro de que esto no se puede combinar en el hilo que hace la E / S?
Si de alguna manera es factible, escribiría un bucle que lea de ambos flujos alternativamente, lo procese / combine y lo escriba en el socket de salida:
while True:
serial_data_in = serial_in.read()
socket_data_in = socket_in.read()
socket_out.write(combine(serial_data_in, socket_data_in))
Tal vez sea necesario modificar los tiempos de espera de los read () s, para evitar perder datos en uno si no hubiera datos entrantes en el otro.
Si eso no funciona , todavía mantendría la menor cantidad de hilos posible. Por ejemplo, puede usar un hilo para la lectura (como arriba) y usar una cola para comunicarse con un hilo que procesa y escribe en el conector de salida:
q = queue.Queue()
def worker_1:
while True:
serial_data_in = serial_in.read()
socket_data_in = socket_in.read()
q.put((serial_data_in, socket_data_in))
def worker_2:
while True:
(serial_data_in, socket_data_in) = q.get()
socket_out.write(combine(serial_data_in, socket_data_in))
q.task_done()
Las colas eliminan la complejidad de sincronización de nivel inferior de bloquear objetos.
Creo que usar select es muy sencillo. Le dice qué socket tiene datos (o EOF) para leer.
En realidad, se ha hecho una pregunta similar antes: Python: el servidor escucha desde dos sockets UDP
Tenga en cuenta que selectse garantiza que solo una lectura de un socket devuelto por no se bloqueará. Verifique nuevamente antes de continuar leyendo. Eso significa que si está leyendo un flujo de datos, lea en un búfer hasta que reciba una línea completa u otra unidad de datos que pueda procesarse.
Su pregunta difiere de la vinculada, porque necesita leer desde la red y una interfaz serial. Linux no tiene ningún problema con él, se puede usar cualquier descriptor de archivo select. Sin embargo, en Windows, solo se pueden usar sockets select. No trabajo con Windows, pero parece que necesitará un hilo dedicado para leer la línea serial.
Puedo sugerir el enfoque utilizado aquí: https://stackoverflow.com/a/641488/4895189. Si tiene una estructura para los datos que recibe a través del conector y la serie, puede escribir esas estructuras con marcas de tiempo en objetos de tubería individuales.
Preferiría multiprocesamiento en lugar de subprocesos de mi experiencia. He usado pyserial para leer y escribir para UART, en el que el hilo principal se usó para escribir y un hilo separado para leer. Por razones que no pude averiguar, perdí fotogramas tanto en la entrada como en la salida si escribía datos sin agregar un retraso bastante grande (~ 1000 ms) entre las llamadas de escritura secuenciales. En general, encuentro que el uso de pyserial con Threading de Python tiene un comportamiento extraño. Actualmente, no estoy seguro si se debe a la implementación de pyserial o al GIL de Python.
Dicho esto, creo que puede usar la siguiente estructura para su configuración según la respuesta que vinculé anteriormente:
Proceso secundario 1: leer datos de Socket y escribir en Pipe con la marca de tiempo
Proceso secundario 2: Leer datos usando pyserial y escribir en Pipe con la marca de tiempo
Proceso principal: Realice la selección en ambos objetos de tubería en un intervalo de su elección, combine los flujos y transmitir a la toma de salida.