So befreien Sie meine Daten von einem Node.js-Stream

Sep 04 2020

Ich habe jetzt eine Weile mit Java Script APIs gearbeitet, aber dies ist das erste Mal, dass ich versucht habe, aus einem aktiven Stream zu testen, der niemals ausgegeben wird 'done'. Mein Ziel ist es, eine festgelegte Anzahl von Proben pro Stunde aus dem Stream zu erhalten. Der Stream verbindet und überträgt viele Informationen, aber ich konnte die zurückgegebenen Daten nicht in ein Format bringen, in dem ich sie weiter verarbeiten kann (wie ich es in einem Data Science-Workflow kenne).

Es fühlt sich so an, als hätte ich seit Tagen auf die Dokumente gestarrt und festgestellt, dass die einfachsten Beispiele den lesbaren Stream in eine Datei auf dem Server leiten. Dies scheint für meine Anwendung ineffizient zu sein. (Sie müssen es in eine Datei schreiben, müssen es jedoch erneut einlesen, um weitere Verarbeitungsschritte auszuführen, bevor Sie es zum Rendern über die Abruf-API an den Browser senden oder zur Langzeitspeicherung und Tiefenanalyse an die MongoDB des Projekts senden. Ich bin mir ziemlich sicher, dass es eine Möglichkeit gibt, JSON als constoder festzulegen, varund ich bin einfach nicht damit vertraut.

Wie bekomme ich meine Daten in die savedJava Script-Variable? Was muss ich ändern oder zu meinem Code hinzufügen, um den zurückgegebenen JSON weiter bearbeiten und verarbeiten zu können?

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

Bearbeitet 2020-09-07

Dies ist ein Beispiel für die Nutzlast der Antwort:

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

Antworten

1 jorgenkg Sep 09 2020 at 12:56

Drei Schritte zur Bewältigung der Herausforderung:

  1. Die Daten müssen als gestreamter HTTP-Antworttext abgerufen werden
  2. Der Antwortstrom muss von einem JSON-Parser analysiert werden, da Daten aus der Antwort gestreamt werden
  3. Der Stream wird beendet, nachdem 20 Elemente vom JSON-Parser analysiert wurden

Der Beispielcode aus dem OP zeigt bereits, wie (1) zu lösen ist.

Es gibt eine Auswahl von Bibliotheken, mit denen Sie einen Stream von JSON-Daten im Handumdrehen analysieren können, um sie zu lösen (2). Meine persönliche Präferenz ist, stream-jsondass nur eine einzige Codezeile in unserer Pipeline erforderlich ist.

Schließlich erfordert (3), dass der Code den eingehenden Stream beendet, bevor er abgeschlossen ist. Dies führt dazu, dass nodejs einen ERR_STREAM_PREMATURE_CLOSEFehler auslöst , der von einer gezielten catch-Anweisung behandelt werden kann.

Das Kombinieren dieser Schritte wird so etwas wie das folgende ausführbare POC. Ich habe kein Twitter-API-Token, aber ich denke, das wird funktionieren:

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