# Making Logstash and/or Filebeats process one entry at a time

**URL:** <https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822>\
**Category:** Logstash\
**Created:** [October 12, 2020, 11:35pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822 "2020-10-12T23:35:07Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 12, 2020, 11:35pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/1 "2020-10-12T23:35:07Z")

</div>

I currently have a setup where Filebeats is tracking a folder (in which there is only 1 file), and sends entries to Logstash. In my logstash.yml, I have pipeline.workers set to 1. Here is a snippet of my logstash.conf, from the "filter" section.

```auto
mutate {
	add_field => {"Entry_time" => "%{[Entry_time_array][0]}"}
	add_field => {"Data_Added_Time" => "%{@timestamp}"}
}

date {
	match => ["Entry_time", "MM/dd/yyyy HH:mm:ss:SSS"]
	target => "@timestamp"
}

grok {
	match => { "Entry_msg" => "TIMING: 'Load Plan Complete', Plan ID '%{NUMBER:taskID}'" }
	add_tag => ["loadplan"]
}

ruby {
	init => "@taskID = -1"
	code => "
		if !event.get('taskID').nil?
			@taskID = event.get('taskID')
		else
			event.set('taskID', @taskID)
		end
	"
}
grok {
	match => {"Entry_msg" => "%{GREEDYDATA:logMessage}"}
}
grok {
	match => { "Entry_msg" => "UI::Received On Treatment Delivery Completed Event" }
	add_tag => ["treatmentcomplete"]
}
grok {
	match => { "Entry_msg" => "Successful Open Delivery Session without Tracking" }
	add_tag => ["treatmentstart"]
}

ruby{
	code => "
		event.set('currentEpochSetTime', Time.now)
		event.set('currentEpoch', event.get('@timestamp').to_i)
	"
}

if [logMessage] =~ "TIMING: 'Load Plan Complete'.*"{
	aggregate{
		aggregate_maps_path => "./aggregate_maps"
		task_id => "%{taskID}"
		code => "
			map['onTableStart'] ||= 0
			map['onTableStart'] = event.get('currentEpoch')
			event.set('onTableStartIs', map['onTableStart'])
			event.set('mapUpdateTime', Time.now)
		"
	}
}

if [logMessage] =~ "Successful Open Delivery Session.*"{
	aggregate{
		task_id => "%{taskID}"
		code => "
			map['treatmentStart'] ||= 0
			map['treatmentStart'] = event.get('currentEpoch')
			event.set('treatmentStartIs', map['treatmentStart'])
			event.set('mapUpdateTime', Time.now)
		"
	}
}

if [logMessage] =~ ".*Received On Treatment Delivery Completed.*" {
	aggregate {
		task_id => "%{taskID}"
		code => "
			event.set('onTableStartIs', map['onTableStart'])
  			event.set('treatmentStartIs', map['treatmentStart'])
			event.set('onTableTime', event.get('currentEpoch') - map['onTableStart'])
  			event.set('treatmentTime', event.get('currentEpoch') - map['treatmentStart'])
			event.set('calculationTime', Time.now)
		"
		map_action => "update"
		end_of_task => true
	}
} else {
	ruby{
		code => "event.set('isNotCompletedEvent', 50)"
	}
}

```

The mutate and the date filters at the top (in addition to some previous filters that are not shown) just make it so that "@timestamp" holds the time specified by the entry (which can be whenever), and "Data\_added\_time" holds the time that the entry was processed by logstash.

For the rest of the config file:

I'm expecting one of three entries to come from filebeats. They are:

1. Load Plan Complete
2. Successful Open Delivery Session
3. Received Treatment Delivery Completed Event

These three entries will always arrive in this order. My goal is to calculate the difference in their @timestamps.

To accomplish this, I am using aggregate maps to track values across multiple entries. The idea is that when I read one of the first two entries, I update the map with a "start time", and when i read the last entry, I use the value in the map to calculate the time difference.

The problem is that I need to guarantee that the map will be updated for messages 1 and 2 BEFORE the calculation happens for message 3. Apparently, this is not guaranteed. I added a couple debug entries to my if-statements. For the first two messages, I added "mapUpdateTime", which marks when the map is updated. For the third message, I added "calculationTime", which marks when the time-difference calculation occured.

So far, it seems guaranteed that Data\_added\_time is always in the correct order. Meaning the log entries are being processed in the correct order. However, sometimes calculationTime will happen before mapUpdateTime, even though calculationTime is a part of the third message, and therefore should be occuring AFTER mapUpdateTime.

My guess so far is that there's still some multithreading going on, but as I mentioned before, I've already set pipeline.workers to 1. I've also tried using the elapsed filter, but that did not work either.

I'm at a loss on what to do, so any help at all would be greatly appreciated. Thanks!

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [October 13, 2020, 1:39am UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/2 "2020-10-13T01:39:54Z")

</div>

> [@jonVR](#):
>
> The problem is that I need to guarantee that the map will be updated for messages 1 and 2 BEFORE the calculation happens for message 3.

I have not read the whole post in detail, so this might be rubbish (it is late for me), but by default the logstash pipeline processes events in batches. 125 events go through the first filter, then 125 events go through the second, etc. If you set pipeline.batch.size to 1 then it might help the function of this particular pipeline, but it might also ruin the scalability of other pipelines. Be careful about setting this without measuring throughput if you have any concerns around scalability.

---

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 13, 2020, 4:26pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/3 "2020-10-13T16:26:34Z")

</div>

Oh, yeah, I do also have the batch size set to 1. It doesn't seem to have any effect on my problem.

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [October 13, 2020, 4:31pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/4 "2020-10-13T16:31:07Z")

</div>

What version of logstash?

---

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 13, 2020, 5:42pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/5 "2020-10-13T17:42:18Z")

</div>

Logstash version is 7.9.0

Filebeats and Elasticsearch version are also 7.9.0

And also Kibana: 7.9.0

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [October 13, 2020, 5:51pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/6 "2020-10-13T17:51:49Z")

</div>

OK, so pipeline.ordered should be auto, and since pipeline.workers is 1 the pipeline should be preserving order. I cannot think of anything else that would result in events getting re-ordered. Can you try setting java\_execution to false, just in case?

---

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 13, 2020, 6:59pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/7 "2020-10-13T18:59:31Z")

</div>

Pipeline.ordered was set to true, but I now have it set to auto. java\_execution was already set to false.

It didn't change anything.

But I want to clarify that the issue is not that the events are getting re-ordered. The ordering is preserved. The issue is that a previous event is not being processed to completion before the next event begins processing.

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [October 13, 2020, 7:24pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/8 "2020-10-13T19:24:35Z")

</div>

> [@jonVR](#):
>
> The ordering is preserved. The issue is that a previous event is not being processed to completion before the next event begins processing.

How can you tell?

---

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 13, 2020, 9:30pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/9 "2020-10-13T21:30:12Z")

</div>

I can tell because mapUpdateTime for the second event happens after Data\_added\_time for the third event. If the second event were processed to completion, then mapUpdateTime for the second event should happen BEFORE Data\_added\_time for the third event.

---

<div class="post-metadata">

**Author:** ![jonVR](https://avatars.discourse-cdn.com/v4/letter/j/7c8e57/32.png) [@jonVR](https://discuss.elastic.co/u/jonVR)\
**Post date:** [October 16, 2020, 10:04pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/10 "2020-10-16T22:04:18Z")

</div>

I have not found a solution to this issue yet. Any additional help would be greatly appreciated.

---

<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:** [November 13, 2020, 10:04pm UTC](https://discuss.elastic.co/t/making-logstash-and-or-filebeats-process-one-entry-at-a-time/251822/11 "2020-11-13T22:04:31Z")

</div>

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