Cómo liberar mis datos de una secuencia de Node.js

Sep 04 2020

He trabajado con las API de Java Script durante un tiempo, pero esta es la primera vez que trato de realizar una muestra de una secuencia activa que nunca se emitirá 'done'. Mi objetivo es obtener una cantidad determinada de muestras de la transmisión por hora. La transmisión está conectando y transmitiendo mucha información, pero no he podido obtener los datos devueltos en un formato en el que pueda procesarlos más (como estoy familiarizado en un flujo de trabajo de ciencia de datos).

Parece que he estado mirando los documentos durante días y noté que la mayoría de los ejemplos sencillos canalizan la transmisión legible en un archivo en el servidor. Esto parece ineficaz para mi aplicación. (Para tener que escribirlo en un archivo, solo para leerlo nuevamente para hacer más procesamiento en él antes de enviarlo al navegador para renderizarlo a través de la API de recuperación o enviarlo al mongoDB del proyecto para un almacenamiento a largo plazo y un análisis profundo. Estoy bastante seguro de que hay una forma de configurar el JSON como consto vary simplemente no estoy familiarizado con él.

¿Cómo introduzco mis datos en la savedvariable Java Script? ¿Qué necesito cambiar o agregar a mi código para poder seguir manipulando y procesando el JSON devuelto?

const needle = require('needle');

const token = process.env.BEARER_TOKEN;
const streamURL = 'https://api.twitter.com/2/tweets/sample/stream';

function streamConnect() {
    const options = {
        timeout: 2000,
    };

    const stream = needle.get(
        streamURL,
        {
            headers: {
                Authorization: `Bearer ${token}`,
            },
        },
        options
    );

    stream
        .on('data', (data) => {
            try {
                const json = JSON.parse(data);
                // console.log(json);
            } catch (e) {
                // Keep alive signal received. Do nothing.
            }
        })
        .on('error', (error) => {
            if (error.code === 'ETIMEDOUT') {
                stream.emit('timeout');
            }
        });

    return stream;
}

function getTweetSample() {
    const s = streamConnect();
    const chunks = [];
    s.on('readable', () => {
        let chunk;
        while (null !== (chunk = s.read())) {
            chunks.push(chunk);
        }
    });
    setInterval(() => {
        s.destroy();
    }, 3000);
    return chunks;
}

const saved = API.getTweetSample();
console.log('saved: ', saved);

// Above returns
// "saved: []"

// Expecting 
// "saved:
{
{
  data: {
    id: '1301578967443337***',
    text: 'See bones too so sure your weight perfect!'
  }
}
{
  data: {
    id: '1301578980001230***
    text: 'Vcs perderam a Dona Maria, ela percebeu q precisa trabalhar e crescer na vida, percebeu q paga 40% de imposto no consumo enquanto políticos q dizem lutar por ela, estão usufruindo dos direitos q ela nunca vai ter 👍 Trabalho escravo é ter q trabalhar pra vcs'
  }
}
...... // 20 examples
}"

Editado 2020-09-07

Esta es una muestra de la carga útil de la respuesta:

PassThrough {
  _readableState: ReadableState {
    objectMode: false,
    highWaterMark: 16384,
    buffer: BufferList { head: null, tail: null, length: 0 },
    length: 0,
    pipes: null,
    pipesCount: 0,
    flowing: true,
    ended: false,
    endEmitted: false,
    reading: false,
    sync: false,
    ....
}

Respuestas

1 jorgenkg Sep 09 2020 at 12:56

Tres pasos para afrontar el desafío:

  1. Los datos deben obtenerse como un cuerpo de respuesta HTTP transmitido
  2. El flujo de respuesta debe ser analizado por un analizador JSON, ya que los datos se transmiten desde la respuesta.
  3. La secuencia terminará después de que el analizador JSON haya analizado 20 elementos.

El código de ejemplo del OP ya ilustra cómo resolver (1).

Existe una selección de bibliotecas para analizar un flujo de datos JSON sobre la marcha para resolverlos (2). Mi preferencia personal es stream-jsonque solo requiere una sola línea de código en nuestra canalización.

Finalmente, (3) requerirá que el código termine la transmisión entrante antes de que se complete. Esto hará que nodejs arroje un ERR_STREAM_PREMATURE_CLOSEerror, que puede ser manejado por una declaración de captura dirigida.

La combinación de estos pasos se convertirá en algo parecido al siguiente POC ejecutable. No tengo un token de API de Twitter, pero creo que esto funcionará:

const stream = require('stream');
const util = require('util');
const got = require('got');
const StreamValues = require("stream-json/streamers/StreamValues.js");

(async () => {
  const token = "<YOUR API TOKEN>";

  const dataStream = got.stream('https://api.twitter.com/2/tweets/sample/stream', {
    headers: { "Authorization": `Bearer ${token}` },
  });

  // This array will by filled by JSON parsed objects from the HTTP response
  const dataPoints = [];
    
  await util.promisify(stream.pipeline)(
    // This readable stream [dataStream] will emit the incoming HTTP body as string data
    dataStream,
    // The string data is then JSON parsed on the fly by [stream-json]
    StreamValues.withParser(),
    // Finally, we iterate over the the JSON objects and push them to the [dataPoints] array.
    async function(source){
      for await (const parsedObject of source){
        dataPoints.push( parsedObject.value );

        if( dataPoints.length === 20 ){
          // When we reach 20 data points, the stream is forcefully terminated
          dataStream.destroy();
          return;
        }
      }
    }
  )
    // Prematurely terminating the stream will cause nodejs to emit a [ERR_STREAM_PREMATURE_CLOSE] 
    // error. If it is OK to return more than 20 elements, you could try to remove the 
    // [return] statement on L28;
    .catch(error => (error.code !== "ERR_STREAM_PREMATURE_CLOSE" && Promise.reject(error)));
}())
  .catch(console.error);