# Logstash pipeline filter

**URL:** <https://discuss.elastic.co/t/logstash-pipeline-filter/294018>\
**Category:** Logstash\
**Created:** [January 11, 2022, 10:28am UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018 "2022-01-11T10:28:07Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![max\_cyril](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/max_cyril/32/100067_2.png) [@max\_cyril](https://discuss.elastic.co/u/max_cyril)\
**Post date:** [January 11, 2022, 10:28am UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/1 "2022-01-11T10:28:07Z")

</div>

Hi,  
I am new to ELK and struggling to write my first logstash pipeline.  
Can anyone help me to write the filter section?  
Thanks in advance.

i want to filter and only output the maximum of completion per users

 ![image (2)](https://us1.discourse-cdn.com/elastic/original/3X/f/c/fc12410b77b2707e0f043ba01305a733ae29d129.png)

for instance :  
user completion  
u2 20  
...  
u48 100

that is my trial whitout succeed.

```auto
input
{elasticsearch
{hosts => "...."
user => "..."
password => "..."
index => "index1"
codec =>"json"
docinfo => true
}}

filter {
aggregate {
task_id => "%{users}"
code => "map['completion'] = event.get('completion') ;
event.cancel if (map['completion']) != map['completion'].max()"
map_action => "create" }

}output
{elasticsearch
{hosts => "..."
user => "..."
password => "..."
index => "index2"
document_type =>"%{[@metadata][_type]}"
document_id =>"%{[@metadata][_id]}"
}}

```

can someone helps me please, thank you!

---

<div class="post-metadata">

**Author:** ![Tomo\_M](https://avatars.discourse-cdn.com/v4/letter/t/848f3c/32.png) [@Tomo\_M](https://discuss.elastic.co/u/Tomo_M)\
**Post date:** [January 11, 2022, 4:12pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/2 "2022-01-11T16:12:47Z")

</div>

Hi,

What is the reason you use logstash?  
If the input and the output is on the same cluster, Transform may be a good option.

you can aggregate data and put into another destination index periodically/continously.

> **[Transforming data | Elasticsearch Guide \[7.16\] | Elastic](https://www.elastic.co/guide/en/elasticsearch/reference/current/transforms.html)**

---

<div class="post-metadata">

**Author:** ![max\_cyril](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/max_cyril/32/100067_2.png) [@max\_cyril](https://discuss.elastic.co/u/max_cyril)\
**Post date:** [January 11, 2022, 4:19pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/3 "2022-01-11T16:19:44Z")

</div>

thank you @Tomo_M ,  
i have several clusters on which i installed Elasticsearch

---

<div class="post-metadata">

**Author:** ![Tomo\_M](https://avatars.discourse-cdn.com/v4/letter/t/848f3c/32.png) [@Tomo\_M](https://discuss.elastic.co/u/Tomo_M)\
**Post date:** [January 11, 2022, 4:54pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/4 "2022-01-11T16:54:01Z")

</div>

then, how about use query in the input and aggregate data without aggregation filter, which is more efficient and you don't have to reinvent the wheel about max aggregation.

I suppose `map['completion']` is a value and not compatible with `max()` function. And, your first line `map['completion'] = event.get('completion')` will update map['completion'] on every event even if `event.get('completion')` is less than `map['completion']`.

---

<div class="post-metadata">

**Author:** ![max\_cyril](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/max_cyril/32/100067_2.png) [@max\_cyril](https://discuss.elastic.co/u/max_cyril)\
**Post date:** [January 11, 2022, 5:14pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/5 "2022-01-11T17:14:00Z")

</div>

yes not compatible with max() function , i got an _non define error function or method_... , i dont master very well my lines of code , i wanted to have a maximum completion value of every user,  
i was thinking first of all must create a map , and then find the maximum value for completion...

---

<div class="post-metadata">

**Author:** ![Tomo\_M](https://avatars.discourse-cdn.com/v4/letter/t/848f3c/32.png) [@Tomo\_M](https://discuss.elastic.co/u/Tomo_M)\
**Post date:** [January 11, 2022, 5:33pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/6 "2022-01-11T17:33:38Z")

</div>

that strategy seems to have some problems because logstash run the `code` event by event.  
logstash can't detect which is the last event for that specific user. `map['completion']` can only keep Array or Value, not both.

I haven't tried it, but how about the following

```auto
event.cancel if !(event.get('completion'));
map['completion'] ||= event.get('completion');
map['completion'] = [map['completion'], event.get('completion')].max;

```

But I strongly recommend to

- use aggregation query in logstash input or
- use aggregation transform in source cluster, and use logstash to simply extract and update.

---

<div class="post-metadata">

**Author:** ![max\_cyril](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/max_cyril/32/100067_2.png) [@max\_cyril](https://discuss.elastic.co/u/max_cyril)\
**Post date:** [January 12, 2022, 11:28am UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/7 "2022-01-12T11:28:26Z")

</div>

```auto
event.cancel if !(event.get('completion'));
map['completion'] ||= event.get('completion');
map['completion'] = [map['completion'], event.get('completion')].max;

```

thank you @Tomo_M ,  
since it gives a result whithout error , the output is not what was expected ...  
**return for each user, the line of higher completion**.  
i've been struggle with it since 2 weeks

---

<div class="post-metadata">

**Author:** ![Tomo\_M](https://avatars.discourse-cdn.com/v4/letter/t/848f3c/32.png) [@Tomo\_M](https://discuss.elastic.co/u/Tomo_M)\
**Post date:** [January 12, 2022, 6:36pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/8 "2022-01-12T18:36:47Z")

</div>

I had to drop the events.

```auto
input{
  elasticsearch{
    hosts => "localhost:9200"
    index => ""
    user => ""
    password => ""
    codec =>"json"
    docinfo => true
    schedule => "* * * * *"
  }
}
filter {
  aggregate {
    task_id => "%{users}"
    code => "event.cancel if !(event.get('completion'));
    map['completion'] ||= event.get('completion');
    map['completion'] = [map['completion'], event.get('completion')].max;
    event.cancel"
    push_map_as_event_on_timeout => true
    timeout=>1
  }
  if !([completion]) {drop{}}
}
output {
  stdout { codec => rubydebug }
}

```

---

<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:** [February 9, 2022, 6:37pm UTC](https://discuss.elastic.co/t/logstash-pipeline-filter/294018/9 "2022-02-09T18:37:05Z")

</div>

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