diff --git a/package-lock.json b/package-lock.json index b5796d64..861c8471 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "aws-lambda-stream", - "version": "1.1.8", + "version": "1.1.9", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "aws-lambda-stream", - "version": "1.1.8", + "version": "1.1.9", "license": "MIT", "dependencies": { "object-sizeof": "^2.6.0" diff --git a/package.json b/package.json index ac1b2630..cfd7f467 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "aws-lambda-stream", - "version": "1.1.8", + "version": "1.1.9", "description": "Create stream processors with AWS Lambda functions.", "keywords": [ "aws", diff --git a/src/sinks/sqs.js b/src/sinks/sqs.js index 9cf458e5..a7cd9300 100644 --- a/src/sinks/sqs.js +++ b/src/sinks/sqs.js @@ -2,7 +2,7 @@ import _ from 'highland'; import Connector from '../connectors/sqs'; -import { toBatchUow, unBatchUow } from '../utils/batch'; +import { batchWithPayloadSizeOrCount, toBatchUow, unBatchUow } from '../utils/batch'; import { ratelimit } from '../utils/ratelimit'; import { rejectWithFault } from '../utils/faults'; import { debug as d } from '../utils/print'; @@ -13,6 +13,7 @@ export const sendToSqs = ({ // eslint-disable-line import/prefer-default-export queueUrl = process.env.QUEUE_URL, messageField = 'message', batchSize = Number(process.env.SQS_BATCH_SIZE) || Number(process.env.BATCH_SIZE) || 10, + maxPayloadSize = Number(process.env.SQS_MAX_PAYLOAD_SIZE) || Number(process.env.MAX_PAYLOAD_SIZE) || 1024 * 1024, parallel = Number(process.env.SQS_PARALLEL) || Number(process.env.PARALLEL) || 8, step = 'send', ...opt @@ -46,7 +47,12 @@ export const sendToSqs = ({ // eslint-disable-line import/prefer-default-export return (s) => s .through(ratelimit(opt)) - .batch(batchSize) + .consume(batchWithPayloadSizeOrCount({ + batchSize, + maxPayloadSize, + payloadField: messageField, + ...opt, + })) .map(toBatchUow) .map(toInputParams) diff --git a/src/utils/batch.js b/src/utils/batch.js index ffa28fc6..050fdd34 100644 --- a/src/utils/batch.js +++ b/src/utils/batch.js @@ -1,5 +1,6 @@ import _ from 'highland'; import isFunction from 'lodash/isFunction'; +import { get } from 'lodash'; import { toClaimcheckEvent, toPutClaimcheckRequest } from '../sinks/claimcheck'; // used after highland batch step @@ -48,6 +49,9 @@ export const compact = (rule) => { }))); }; +/** + * Batch EB request entries by size to avoid writing a batch that's too large to EB. + */ export const batchWithSize = ({ claimCheckBucketName = process.env.CLAIMCHECK_BUCKET_NAME, putClaimcheckRequest = 'putClaimcheckRequest', @@ -116,3 +120,58 @@ const logMetrics = (batch, sizes, opt) => { batch[0].metrics?.gauge('publish|stream.pipeline.eventSize.bytes', sizes); } }; + +/** + * Batches by aggregate payload size with a cap on payload count. + */ +export const batchWithPayloadSizeOrCount = ({ + batchSize, + maxPayloadSize, + payloadField, + ...opt +}) => { + let batched = []; + let sizes = []; + + return (err, x, push, next) => { + /* istanbul ignore if */ + if (err) { + push(err); + next(); + } else if (x === _.nil) { + if (batched.length > 0) { + logMetrics(batched, sizes, opt); + push(null, batched); + } + + push(null, _.nil); + } else { + if (!get(x, payloadField)) { + push(null, [x]); + } else { + const size = Buffer.byteLength(JSON.stringify(get(x, payloadField))); + if (size > maxPayloadSize) { + logMetrics([x], [size], opt); + const error = new Error(`Payload size: ${size}, exceeded max: ${maxPayloadSize}`); + error.uow = x; + push(error); + next(); + return; + } + + const totalSize = sizes.reduce((a, c) => a + c, size); + if (totalSize <= maxPayloadSize && batched.length + 1 <= batchSize) { + batched.push(x); + sizes.push(size); + } else { + logMetrics(batched, sizes, opt); + push(null, batched); + batched = [x]; + sizes = [size]; + } + } + + next(); + } + }; +}; diff --git a/test/unit/sinks/sqs.test.js b/test/unit/sinks/sqs.test.js index c088cb21..96fd5279 100644 --- a/test/unit/sinks/sqs.test.js +++ b/test/unit/sinks/sqs.test.js @@ -18,6 +18,11 @@ describe('sinks/sqs.js', () => { Id: '1', MessageBody: JSON.stringify({ f1: 'v1' }), }, + }, { + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, }]; _(uows) @@ -26,7 +31,68 @@ describe('sinks/sqs.js', () => { .tap((collected) => { // console.log(JSON.stringify(collected, null, 2)); - expect(collected.length).to.equal(1); + expect(collected.length).to.equal(2); + expect(collected[0]).to.deep.equal({ + message: { + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, + inputParams: { + Entries: [{ + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }], + }, + sendMessageBatchResponse: {}, + }); + + expect(collected[1]).to.deep.equal({ + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, + inputParams: { + Entries: [{ + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }], + }, + sendMessageBatchResponse: {}, + }); + }) + .done(done); + }); + + it('should split a batch due to batch size', (done) => { + sinon.stub(Connector.prototype, 'sendMessageBatch').resolves({}); + + const uows = [{ + message: { + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, + }, { + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, + }]; + + _(uows) + .through(sendToSqs({ + batchSize: 1, + })) + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + + expect(collected.length).to.equal(2); expect(collected[0]).to.deep.equal({ message: { Id: '1', @@ -40,6 +106,73 @@ describe('sinks/sqs.js', () => { }, sendMessageBatchResponse: {}, }); + expect(collected[1]).to.deep.equal({ + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, + inputParams: { + Entries: [{ + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }], + }, + sendMessageBatchResponse: {}, + }); + }) + .done(done); + }); + + it('should split a batch due to payload size', (done) => { + sinon.stub(Connector.prototype, 'sendMessageBatch').resolves({}); + + const uows = [{ + message: { + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, + }, { + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, + }]; + + _(uows) + .through(sendToSqs({ + maxPayloadSize: 50, + })) + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + + expect(collected.length).to.equal(2); + expect(collected[0]).to.deep.equal({ + message: { + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }, + inputParams: { + Entries: [{ + Id: '1', + MessageBody: JSON.stringify({ f1: 'v1' }), + }], + }, + sendMessageBatchResponse: {}, + }); + expect(collected[1]).to.deep.equal({ + message: { + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }, + inputParams: { + Entries: [{ + Id: '2', + MessageBody: JSON.stringify({ f1: 'v2' }), + }], + }, + sendMessageBatchResponse: {}, + }); }) .done(done); }); diff --git a/test/unit/utils/batch.test.js b/test/unit/utils/batch.test.js index f3ffee63..23e750c4 100644 --- a/test/unit/utils/batch.test.js +++ b/test/unit/utils/batch.test.js @@ -5,6 +5,7 @@ import _ from 'highland'; import { toBatchUow, unBatchUow, group, batchWithSize, compact, + batchWithPayloadSizeOrCount, } from '../../../src/utils'; describe('utils/batch.js', () => { @@ -439,4 +440,218 @@ describe('utils/batch.js', () => { }) .done(done); }); + + describe('batchWithPayloadSizeOrCount', () => { + it('should batch on size', (done) => { + const uows = [ + { + message: { // size = 19 + id: 'xxxxxxxxxx', + }, + }, + { + message: { // size = 29 + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { // size = 39 + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { // size = 10 + id: 'x', + }, + }, + ]; + + _(uows) + .consume(batchWithPayloadSizeOrCount({ + batchSize: 999, + maxPayloadSize: 50, + payloadField: 'message', + })) + + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + + expect(collected.length).to.equal(2); + expect(collected[0]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxx', + }, + }, + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + ]); + expect(collected[1]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { + id: 'x', + }, + }, + ]); + }) + .done(done); + }); + + it('should batch on count', (done) => { + const uows = [ + { + message: { // size = 19 + id: 'xxxxxxxxxx', + }, + }, + { + message: { // size = 29 + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { // size = 39 + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + ]; + + _(uows) + .consume(batchWithPayloadSizeOrCount({ + batchSize: 2, + maxPayloadSize: 999, + payloadField: 'message', + // metricsEnabled: true, + // debug: (msg, v) => console.log(msg, v), + })) + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + + expect(collected.length).to.equal(2); + expect(collected[0]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxx', + }, + }, + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + ]); + expect(collected[1]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + ]); + }) + .done(done); + }); + + it('should handle oversized requests', (done) => { + const spy = sinon.spy(); + const uows = [ + { + message: { // size = 19 + id: 'xxxxxxxxxx', + }, + }, + { + message: { // size = 39 + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { // size = 29 + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + ]; + + _(uows) + .consume(batchWithPayloadSizeOrCount({ + batchSize: 2, + maxPayloadSize: 30, + payloadField: 'message', + })) + .errors(spy) + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + expect(collected.length).to.equal(2); + expect(spy).to.have.been.calledWith; // (Error('Request size: 39, exceeded max: 30')); + }) + .done(done); + }); + + it('should skip batching nonexist payload fields', (done) => { + const uows = [ + { + message: { // size = 19 + id: 'xxxxxxxxxx', + }, + }, + { + message: { // size = 29 + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + { + message: { // size = 39 + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + ]; + + _(uows) + .consume(batchWithPayloadSizeOrCount({ + batchSize: 2, + maxPayloadSize: 999, + payloadField: 'fake-field', + // metricsEnabled: true, + // debug: (msg, v) => console.log(msg, v), + })) + .collect() + .tap((collected) => { + // console.log(JSON.stringify(collected, null, 2)); + + expect(collected.length).to.equal(3); + expect(collected[0]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxx', + }, + }, + ]); + expect(collected[1]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxx', + }, + }, + ]); + expect(collected[2]).to.deep.equal([ + { + message: { + id: 'xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx', + }, + }, + ]); + }) + .done(done); + }); + }); });