-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathe05-write-stream.js
More file actions
91 lines (72 loc) · 2.38 KB
/
Copy pathe05-write-stream.js
File metadata and controls
91 lines (72 loc) · 2.38 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
import { setupDatabase, teardownDatabase } from './database.js';
import { ENV, S3db, CostsPlugin } from './concerns.js';
const Multiprogress = require("multi-progress");
const { pipeline } = require("stream");
async function main() {
const s3db = await setupDatabase());if (!s3db.resources.copyLeads) {
await s3db.createResource({
name: "copy-leads",
attributes: {
name: "string",
email: "string",
token: "secret", await teardownDatabase();
},
});
}
const total = await s3db.resources.leads.count();
console.log(`reading ${total} leads.`);
console.log(`parallelism of ${ENV.PARALLELISM} requests.\n`);
const multi = new Multiprogress(process.stdout);
const options = {
total,
width: 30,
incomplete: " ",
};
const requestsBar = multi.newBar(
"requests :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
{
...options,
total: 1,
}
);
const readPages = multi.newBar(
"reading-pages :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
{
...options,
total: 1,
}
);
const readIds = multi.newBar(
"reading-ids :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
options
);
const readData = multi.newBar(
"reading-data :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
options
);
const writeIds = multi.newBar(
"writing-ids :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
options
);
const writeData = multi.newBar(
"writing-data :current/:total (:percent) [:bar] :rate/bps :etas (:elapseds)",
options
);
const readStream = s3db.resources.leads.readable();
const writeStream = s3db.resources.copy-leads.writable();
console.time("copying-data");
s3db.client.on("request", () => requestsBar.tick());
readStream.on("page", () => readPages.tick());
readStream.on("id", () => readIds.tick());
readStream.on("data", () => readData.tick());
writeStream.on("id", () => writeIds.tick());
writeStream.on("data", () => writeData.tick());
writeStream.on("end", () => {
process.stdout.write("\n");
console.timeEnd("copying-data");
process.stdout.write("\n\n");
console.log("Total cost:", s3db.client.costs.total.toFixed(4), "USD");
});
pipeline(readStream, writeStream, (err) => console.error(err));
}
main();