Comment libérer mes données d'un flux Node.js

Sep 04 2020

Je travaille avec les API Java Script depuis un certain temps maintenant, mais c'est la première fois que j'essaie d'échantillonner à partir d'un flux actif qui n'émettra jamais 'done'. Mon objectif est d'obtenir un nombre défini d'échantillons du flux par heure. Le flux se connecte et diffuse beaucoup d'informations, mais je n'ai pas été en mesure d'obtenir les données renvoyées dans un format permettant de les traiter ultérieurement (comme je le connais dans un flux de travail de science des données).

J'ai l'impression de regarder la documentation depuis des jours maintenant, et j'ai remarqué que la plupart des exemples simples canalisent le flux lisible dans un fichier sur le serveur. Cela semble inefficace pour mon application. (Pour avoir à l'écrire dans un fichier, seulement pour le relire pour faire plus de traitement dessus avant de l'envoyer au navigateur pour le rendu via l'API de récupération ou de l'envoyer à la mongoDB du projet pour un stockage à long terme et une analyse approfondie. Je suis à peu près sûr qu'il existe un moyen de définir le JSON en tant que constou varet je ne suis tout simplement pas familier avec cela.

Comment obtenir mes données dans la savedvariable Java Script? Que dois-je modifier ou ajouter à mon code pour pouvoir continuer à manipuler et à traiter le JSON renvoyé?

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
}"

Modifié 2020-09-07

Voici un exemple de la charge utile de la réponse:

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,
    ....
}

Réponses

1 jorgenkg Sep 09 2020 at 12:56

Trois étapes pour relever le défi:

  1. Les données doivent être récupérées en tant que corps de réponse HTTP diffusé
  2. Le flux de réponse doit être analysé par un analyseur JSON car les données sont diffusées à partir de la réponse
  3. Le flux doit se terminer après que 20 éléments ont été analysés par l'analyseur JSON

L'exemple de code de l'OP illustre déjà comment résoudre (1).

Il existe une sélection de bibliothèques pour analyser un flux de données JSON à la volée à résoudre (2). Ma préférence personnelle est stream-jsoncar il ne nécessite qu'une seule ligne de code dans notre pipeline.

Enfin, (3) exigera que le code termine le flux entrant avant qu'il ne se termine. Cela entraînera le renvoi d'une ERR_STREAM_PREMATURE_CLOSEerreur par nodejs , qui peut être gérée par une instruction catch ciblée.

La combinaison de ces étapes deviendra quelque chose comme le POC exécutable suivant. Je n'ai pas de jeton d'API Twitter, mais je pense que cela fonctionnera:

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);