Cara membebaskan data saya dari Aliran Node.js
Saya telah bekerja dengan Java Script API untuk sementara waktu sekarang, tetapi ini adalah pertama kalinya saya mencoba mengambil sampel dari aliran aktif yang tidak akan pernah memancarkan 'done'. Tujuan saya adalah mendapatkan sejumlah sampel dari aliran per jam. Aliran menghubungkan dan mengalirkan banyak informasi, tetapi saya belum bisa mendapatkan data yang dikembalikan ke dalam format di mana saya dapat melakukan pemrosesan lebih lanjut padanya (seperti yang saya kenal dalam alur kerja ilmu data).
Rasanya seperti saya telah menatap dokumen selama berhari-hari sekarang, dan melihat sebagian besar contoh langsung menyalurkan aliran yang dapat dibaca ke dalam file di server. Ini sepertinya tidak efisien untuk aplikasi saya. (Harus menulisnya ke file, hanya untuk membacanya lagi untuk melakukan lebih banyak pemrosesan sebelum kemudian mengirimnya ke browser untuk dirender melalui API pengambilan atau mengirimnya ke mongoDB proyek untuk penyimpanan jangka panjang dan analisis mendalam. Saya cukup yakin ada cara untuk menyetel JSON sebagai constatau vardan saya tidak terbiasa dengannya.
Bagaimana cara memasukkan data saya ke dalam savedvariabel Java Script? Apa yang saya perlukan untuk mengubah atau menambahkan kode saya agar dapat terus memanipulasi dan memproses JSON yang dikembalikan?
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
}"
Diedit 2020-09-07
Ini adalah contoh payload dari respon:
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,
....
}
Jawaban
Tiga langkah untuk mengatasi tantangan:
- Data harus diambil sebagai badan respons HTTP yang dialirkan
- Aliran respons harus diurai oleh pengurai JSON karena data dialirkan dari respons
- Aliran akan berhenti setelah 20 elemen telah diurai oleh parser JSON
Contoh kode dari OP sudah menggambarkan cara menyelesaikan (1).
Ada pilihan perpustakaan di luar sana untuk mengurai aliran data JSON dengan cepat untuk dipecahkan (2). Preferensi pribadi saya adalah stream-jsonkarena hanya membutuhkan satu baris kode di jalur pipa kami.
Terakhir, (3) akan membutuhkan kode untuk menghentikan aliran masuk sebelum selesai. Ini akan menyebabkan nodejs melontarkan ERR_STREAM_PREMATURE_CLOSEkesalahan, yang bisa ditangani oleh pernyataan catch yang ditargetkan.
Menggabungkan langkah-langkah ini akan menjadi sesuatu seperti POC yang dapat dieksekusi berikut ini. Saya tidak memiliki token API Twitter, tetapi menurut saya ini akan berfungsi:
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);