I have two cron-jobs running parallel and reuse an ElasticSearch client for importing data to a pelias instance or ES. Currently, I'm using the wof-admin-lookup package to transform the data. And this package uses another package call parallel-transform for parallel transforming.
Here how the parallel transforming is used:
module.exports = function(pipResolver, config) {
if (!pipResolver) {
throw new Error('valid pipResolver required to be passed in as the first parameter');
}
// pelias 'imports.adminLookup' config section
config = config || {};
const pipResolverStream = createPipResolverStream(pipResolver, config);
const end = createPipResolverEnd(pipResolver);
const stream = parallelTransform(config.maxConcurrentReqs || 1, pipResolverStream);
stream.on('end', end);
return stream;
};
Here is my code use wof-admin-lookup:
const stream = recordStream.create(filePath)
.pipe(blacklistStream())
.pipe(adminLookup.create())
.pipe(model.createDocumentMapperStream())
.pipe(peliasDbclient())
Everything is working well, but I'm facing an issue that sometimes happens. If there is no data piped to the adminLookup which means there is no parallel transform, the "end" event will never fire. I plan to modify the parallel-transform package code, so I wonder if there is any way to end a stream after a period of time that has no data piped? I tried to switch the stream mode by using the 'data' event but it is helpless. Thank you.