# Aggregate filter for mixed data lines

**URL:** <https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110>\
**Category:** Logstash\
**Created:** [June 8, 2020, 1:19am UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110 "2020-06-08T01:19:50Z")\
**Posts on this page:** 6\
**Page:** 1

<div class="post-metadata">

**Author:** ![mskadu](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mskadu/32/69520_2.png) [@mskadu](https://discuss.elastic.co/u/mskadu)\
**Post date:** [June 8, 2020, 1:19am UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/1 "2020-06-08T01:19:50Z")

</div>

Hello all!

I am trying to use the Aggregate filter plugin to correlate and combine data from two different CSV file inputs which represent API data calls. The idea is to produce a record showing a combined picture. As you can expect the data may or may not arrive in the right sequence. Here's an example:

/data/incoming/source\_1/\*.csv

> StartTime, AckTime, Operation, RefData1, RefData2, OpSpecificData1  
> 231313232,44343545,Register,ref-data-1a,ref-data-2a,op-specific-data-1  
> 979898999,75758383,Register,ref-data-1b,ref-data-2b,op-specific-data-2  
> 354656466,98554321,Cancel,ref-data-1c,ref-data-2c,op-specific-data-2

/data/incoming/source\_1/\*.csv

> FinishTime,Operation,RefData1, RefData2, FinishSpecificData  
> 67657657575,Cancel,ref-data-1c,ref-data-2c,FinishSpecific-Data-1  
> 68445590877,Register,ref-data-1a,ref-data-2a,FinishSpecific-Data-2  
> 55443444313,Register,ref-data-1a,ref-data-2a,FinishSpecific-Data-2

I have a single pipeline that is receiving both these CSVs and I am able to process and write them as individual records to a single Index. However, the idea is to combine records from the two sources into one record each representing a superset. of Operation related information

Unfortunately, despite several attempts I have been unable to figure out how to achieve this via Aggregate filter plugin. My primary question is whether this is a suitable use of the specific plugin? And if so, any suggestions would be welcome!

At the moment, I have this

```auto
input {
   file {
      path => ['/data/incoming/source_1/*.csv']
      tags => ["source1"]
   }
   file {
      path => ['/data/incoming/source_2/*.csv']
      tags => ["source2"]
   }
   # use the tags to do some source 1 and 2 related massaging, calculations, etc

   aggregate {
         task_id = "%{Operation}_%{RefData1}_%{RefData1}"
         code => "
             map['source_files'] ||= []
             map['source_files'] << {'source_file', event.get('path') }
         "
         push_map_as_event_on_timeout => true
         timeout => 600 #assuming this is the most far apart they will arrive         
   }
  ...
}
output {
    elastic { ...}
}

```

And other such variations. However, I keep getting individual records being written to the Index and am unable to get one combined. Yet again, as you can see from the data set there's no guarantee of the sequencing of records - so I am wondering if the filter is the right tool for the job, to begin with? 🙄

Or is it just me not being able to use it right! 🤭

In either case, any inputs/ comments/ suggestions welcome. Thanks!

Edit: Also cross-posted [to Stackoverflow](https://stackoverflow.com/questions/62262043/using-logstash-aggregate-filter-plugin-to-process-data-which-may-or-may-not-be-s).

---

<div class="post-metadata">

**Author:** ![Rahul\_Kumar4](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rahul_kumar4/32/67369_2.png) [@Rahul\_Kumar4](https://discuss.elastic.co/u/Rahul_Kumar4)\
**Post date:** [June 8, 2020, 3:40pm UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/2 "2020-06-08T15:40:08Z")

</div>

> [@mskadu](#):
>
> ```auto
> file {
> path => ['/data/incoming/source_1/*.csv']
> tags => ["source1"]
> }
> file {
> path => ['/data/incoming/source_2/*.csv']
> tags => ["source2"]
> }
> 
> ```

The two `file` input plugins in your Logstash pipelines run in their own threads and do not share events with each other as explained in the execution model [here.] ([Execution Model | Logstash Reference [8.11] | Elastic](https://www.elastic.co/guide/en/logstash/current/execution-model.html))

To overcome this problem, you could do something like below:

1. In pipeline (lets say A), first index the CSVs say `/data/incoming/source_1/*.csv` in Elasticsearch `indexA`

2. In another pipeline (lets say B), use the `file` input plugin to read the source`/data/incoming/source_2/*.csv` and use the [Elasticsearch filter plugin](https://www.elastic.co/guide/en/logstash/current/plugins-filters-elasticsearch.html) to query the `indexA` to look up the corresponding values and enrich your events and then let this get indexed in `indexB`. `indexB` is what you are looking for.

---

<div class="post-metadata">

**Author:** ![mskadu](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mskadu/32/69520_2.png) [@mskadu](https://discuss.elastic.co/u/mskadu)\
**Post date:** [June 8, 2020, 5:18pm UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/3 "2020-06-08T17:18:41Z")

</div>

That sounds like an idea! Let me give it a shot and come back with results o further question. 😀 [quote="Rahul\_Kumar4, post:2, topic:236110"]  
use the [Elasticsearch filter plugin](https://www.elastic.co/guide/en/logstash/current/plugins-filters-elasticsearch.html) to query the `indexA` to look up the corresponding values and enrich your events and then let this get indexed in `indexB` . `indexB` is what you are looking for.  
[/quote]

---

<div class="post-metadata">

**Author:** ![mskadu](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/mskadu/32/69520_2.png) [@mskadu](https://discuss.elastic.co/u/mskadu)\
**Post date:** [June 8, 2020, 6:20pm UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/4 "2020-06-08T18:20:35Z")

</div>

In the meanwhile, I spotted [this post](https://discuss.elastic.co/t/import-csv-with-different-column-names-to-same-field/223725/4) which allows me to use the Elasticsearch output plugin in upsert mode - which pretty near does what i need. Here's what my output section now looks like

```auto
...
output {

  if "source1" in [tags] {
       elasticsearch { ..} # write to source1 specific index
  }
  else if "source2" in [tags] {
       elasticsearch { ..} # write to source2 specific index
  }
  # and ultimately the combined index
   elasticsearch{
      hosts => ["my-es-host:9200"]
      index => ["my-combined-index"]
      action => "update"
      document_id => "%{Operation}_%{RefData1}_%{RefData1}"
      doc_as_upsert => true
   }
}

```

This gives me two choices - wicked!

---

<div class="post-metadata">

**Author:** ![Rahul\_Kumar4](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rahul_kumar4/32/67369_2.png) [@Rahul\_Kumar4](https://discuss.elastic.co/u/Rahul_Kumar4)\
**Post date:** [June 8, 2020, 7:54pm UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/5 "2020-06-08T19:54:03Z")

</div>

> [@mskadu](#):
>
> In the meanwhile, I spotted [this post](https://discuss.elastic.co/t/import-csv-with-different-column-names-to-same-field/223725/4) which allows me to use the Elasticsearch output plugin in upsert mode - which pretty near does what i need. Here's what my output section now looks like

Yeap, this should work too. You won't have to use the Elasticsearch in the filter section in this case and you won't have the redundant `indexB` in this case as your `indexA` will itself get updated.

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [July 6, 2020, 7:54pm UTC](https://discuss.elastic.co/t/aggregate-filter-for-mixed-data-lines/236110/6 "2020-07-06T19:54:10Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
