Problema: tengo un atributo de objeto .status que se actualiza con un tema de Kafka, luego lo envío a través de websocket. Mi problema es que cada vez que pregunto desde el lado del Cliente (javascript), el Servidor (websockets + asyncio en Python) iniciará un consumidor de Kafka desde el principio.
Pregunta: ¿Es posible tener el bucle for de Kafka ( for msg in consumer: :) actualizando mi objeto custom_obj y enviar su valor .status solo cuando se lo solicite?
Esto es lo que tengo hasta ahora en el lado del servidor:
import asyncio import websockets from kafka import KafkaConsumer import Custom_obj async def test(websocket): consumer = KafkaConsumer( 'kafka-topic', bootstrap_servers=['kafka.server.com:1234'], auto_offset_reset='earliest', #Must start from the beginning to build the object correctly enable_auto_commit=True, ) custom_obj = Custom_obj() for msg in consumer: msg_dec = msg.value.decode() custom_obj.update(msg_dec) await websocket.send(custom_obj.status) async def main(): async with websockets.serve(test, "localhost", 1234): await asyncio.Future() # run forever if __name__ == "__main__": asyncio.run(main())Código del lado del cliente (javascript en el componente Vue):
created() { const ws = new WebSocket('ws://localhost:1234'); ws.onopen = function(e) { ws.send('Got here!') this.connectionStatus = 'Connected.' } ws.onerror = function(e) { ws.close() } ws.onclose = function(e) { this.connectionStatus = 'Disconnected.' } ws.onmessage = (e) => { console.log(1) }Si desea realizar un seguimiento del progreso en el tema de Kafka, deberá usar un parámetro group_id en el constructor.
También es posible que desee consultar aiokafka para obtener soporte asincrónico adicional.