Tengo un punto final de API que transmite una respuesta JSON. Ahora quiero usar los observables RxJs para transmitir los fragmentos de datos de esta manera:

Y suscríbete a lo observable en otras clases. ¿Cómo puedo hacerlo? ¡Gracias!
Lo consumo de esta manera:
http.get({ hostname: port: 80, path: '/api/' + "v4" + '/contests/' + "contest" + '/event-feed', method: 'GET', headers: { 'Content-Type': 'application/json' }, qs: { strict: false, stream: true } }, (res) => { res.on('data', (chunk) => { rawData += chunk; let obj; try { obj = JSON.parse(chunk); console.log(obj) } catch (e) { if (e.constructor.name != "SyntaxError") console.log("[ERROR]: " + e); else console.log("..."); } }); res.on('end', () => { try { console.log("---FIN---\n"); const parsedData = JSON.parse(rawData); console.log(parsedData); } catch (e) { console.error(e.message); } });Y funciona.
En realidad, hay mucha sutileza en lo que está preguntando aquí y cómo interactúa con los Observables, en lo que intentaré profundizar un poco ...
Una forma de hacer esto es crear un Asunto antes de activar el http.get que luego llama al método .next() cuando recibe datos. Devuelves el Subject y luego te suscribes a eso. El desafío aquí es que es posible que la operación http ya se haya completado en el momento en que se suscriba y pierda los fragmentos de datos.
Una forma de evitar perder cosas es usar ReplaySubject en su lugar, que almacenará en búfer y reproducirá algunos valores para los nuevos suscriptores.
La forma más habitual de hacerlo es no activar la llamada http hasta que alguien se suscriba, lo que puede hacer así...
var http = require('https'); var { Observable } = require('rxjs'); function prepareConnection() { return new Observable(subscriber => { http.get({ hostname: port: 80, path: '/api/' + "v4" + '/contests/' + "contest" + '/event-feed', method: 'GET', headers: { 'Content-Type': 'application/json' }, qs: { strict: false, stream: true } }, (res) => { res.on('data', (chunk) => { let obj; try { obj = JSON.parse(chunk); subscriber.next(obj); } catch (e) { subscriber.error(e); } }); res.on('end', () => { subscriber.complete(); }); }); }); } prepareConnection().subscribe({ next(x) { console.log(x); }, error(err) { console.error('something wrong occurred: ' + err); }, complete() { console.log('done'); } }); La función que se pasa al constructor se ejecuta cuando alguien se suscribe al Observable , siendo el suscriptor el primer argumento de esa función a la que puede llamar next() , etc.
La advertencia aquí es que cada nuevo suscriptor activará una nueva llamada http. Ese puede o no ser el comportamiento que desea. Si no es lo que desea, puede usar compartir para compartir una suscripción, por ejemplo.
let connection = prepareConnection().pipe(share()); connection.subscribe({ ... }); connection.subscribe({ ... });