forked from mapbox/dynamodb-replicator
-
Notifications
You must be signed in to change notification settings - Fork 6
/
Copy paths3-backfill.js
88 lines (72 loc) · 2.59 KB
/
s3-backfill.js
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
var Dyno = require('dyno');
var AWS = require('aws-sdk');
var s3 = new AWS.S3({
maxRetries: 1000,
httpOptions: {
timeout: 1000,
agent: require('streambot').agent
}
});
var stream = require('stream');
var queue = require('queue-async');
var crypto = require('crypto');
module.exports = function(config, done) {
var primary = Dyno(config);
if (config.backup)
if (!config.backup.bucket || !config.backup.prefix)
return done(new Error('Must provide a bucket and prefix for incremental backups'));
primary.describeTable(function(err, data) {
if (err) return done(err);
var keys = data.Table.KeySchema.map(function(schema) {
return schema.AttributeName;
});
var count = 0;
var starttime = Date.now();
var writer = new stream.Writable({ objectMode: true, highWaterMark: 1000 });
writer.queue = queue();
writer.queue.awaitAll(function(err) { if (err) done(err); });
writer.pending = 0;
writer._write = function(record, enc, callback) {
if (writer.pending > 1000)
return setImmediate(writer._write.bind(writer), record, enc, callback);
var key = keys.reduce(function(key, k) {
key[k] = record[k];
return key;
}, {});
var id = crypto.createHash('md5')
.update(Dyno.serialize(key))
.digest('hex');
var params = {
Bucket: config.backup.bucket,
Key: [config.backup.prefix, config.table, id].join('/'),
Body: Dyno.serialize(record)
};
writer.drained = false;
writer.pending++;
writer.queue.defer(function(next) {
s3.putObject(params, function(err) {
count++;
process.stdout.write('\r\033[K' + count + ' - ' + (count / ((Date.now() - starttime) / 1000)).toFixed(2) + '/s');
writer.pending--;
if (err) writer.emit('error', err);
next();
});
});
callback();
};
writer.once('error', done);
var end = writer.end.bind(writer);
writer.end = function() {
writer.queue.awaitAll(end);
};
primary.scanStream()
.on('error', next)
.pipe(writer)
.on('error', next)
.on('finish', next);
function next(err) {
if (err) return done(err);
done(null, { count: count });
}
});
};