I have a pipeline from an SQL database to Elasticsearch that looks something like the following:
- Input from the SQL database using logstash-input-jdbc
- Various filtering and mutation of the individual events
- The events are aggregated according to a group_id property using logstash-filter-aggregate
- The aggregate events are output to Elasticsearch using logstash-output-elasticsearch
As it is, the throughput of this pipeline is quite low. I know that this is due to the aggregation step (which performs some relatively heavy processing), and I would like to use several threads/processes in order to improve performance (allowing me to utilize more than one core).
However, the logstash-filter-aggregate plugin does not support multiple filter workers -- presumably because it has no way to guarantee that the events that should be combined into one aggregate events will be processed by the same worker.
My current solution to this is to run several instances of logstash where each instance selects a certain subset of group_ids from the SQL database. However, there is quite a bit of overhead to this. Are there any better ways to use multiple cores with logstash-filter-aggregate?