# Kafka input - resync missing items from topic

**URL:** <https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294>\
**Category:** Logstash\
**Created:** [February 10, 2023, 8:45pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294 "2023-02-10T20:45:27Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![Chris\_Denneen](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chris_denneen/32/1254_2.png) [@Chris\_Denneen](https://discuss.elastic.co/u/Chris_Denneen)\
**Post date:** [February 10, 2023, 8:45pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/1 "2023-02-10T20:45:27Z")

</div>

Ran into issue where we have logstash input reading kafka topics for log events.  
Last night the ES index rolled over and the write alias was lost (not on the new index... so nothing with is\_write\_index = true) therefore logstash failed for 15 hours failing to write to ES.

We fixed the alias and "new" data started coming in.

My question is none of the last 15 hours are showing up or even beginning to resync into the index.  
Is there a way to have the kafka input "resync" to see what it missed to index?  
I obviously do not want to roll back to the beginning of the topic or to duplicate already indexed events but I'd like to get back the last 15 hours for events we are missing.

Thanks

---

<div class="post-metadata">

**Author:** ![leandrojmp](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/leandrojmp/32/107231_2.png) [@leandrojmp](https://discuss.elastic.co/u/leandrojmp)\
**Post date:** [February 11, 2023, 2:30pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/2 "2023-02-11T14:30:59Z")

</div>

Do you have the offset of the last message indexed before your issue and the first one after you fixed?

If you have both offsets you may be able to spin-up another pipeline and configure it to consume from the beginning of your topic, but create a filter in the logstash pipeline to drop anything before and after those offsets.

If you do not have the offsets you may do a similar thing, spin-up another pipeline to consume the topic from the beginning, but send the data to a temporary index, after the date in this index reaches the current data you may stop it.

Then you would need to do a delete\_by\_query to remove the data from before the issue and after the issue is fixed and reindex it in your index.

Those are the two ways I could think of.

---

<div class="post-metadata">

**Author:** ![Chris\_Denneen](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chris_denneen/32/1254_2.png) [@Chris\_Denneen](https://discuss.elastic.co/u/Chris_Denneen)\
**Post date:** [February 11, 2023, 3:35pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/3 "2023-02-11T15:35:50Z")

</div>

Unless it logs it somewhere I'm not sure I would be able to tell where to get this info from.  
Is there any sort of improvement that can be done on the input that it wouldn't mark it in offset as complete until it gets confirmation about "output" being successful?

I mean the logstash config is literally input(kafka) -\> filter -\> output(search). If it doesn't successfully index due to this write issue or disk full or countless other reasons I would want kafka to keep the offset to where the last successful index item was.

I haven't specified anything custom on logstash config but believe it could/can basically buffer itself so while kafka didn't keep the offset those items that failed to send to search should still be in logstash buffer to retry to search? Failing to send shouldn't allow it to be "dropped" from retry logic I wouldn't think.

---

<div class="post-metadata">

**Author:** ![Chris\_Denneen](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chris_denneen/32/1254_2.png) [@Chris\_Denneen](https://discuss.elastic.co/u/Chris_Denneen)\
**Post date:** [February 11, 2023, 4:08pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/4 "2023-02-11T16:08:22Z")

</div>

Ok so it won't help me with this specific event but I found that I can capture and log the kafka offset from the metadata. So in the future I should be able to at least find out the last offset and the "first" when it picks up again.

Do you have any sort of example on the logstash pipeline filter to drop events outside those 2 offsets (so basically I'd have to setup the kafka input to use a temp consumer group, set to beginning so it goes back to the topic start)... then have it drop all events before the last event and drop all events after the "first" event from pickup... and once it's done processing remove that pipeline. (I'd obviously keep it around for one off's like this in the future).

---

<div class="post-metadata">

**Author:** ![leandrojmp](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/leandrojmp/32/107231_2.png) [@leandrojmp](https://discuss.elastic.co/u/leandrojmp)\
**Post date:** [February 11, 2023, 4:53pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/5 "2023-02-11T16:53:41Z")

</div>

> [@Chris\_Denneen](#):
>
> If it doesn't successfully index due to this write issue or disk full or countless other reasons I would want kafka to keep the offset to where the last successful index item was.

If I'm not wrong the input does not have an end-to-end ack, so it does not wait the event to be indexed in elasticsearch to commit the offset, it commits the offset periodically after they reach the queue state of the pipeline (memory or persistent queue).

Every logstash pipeline has the following flow:

input -\> queue -\> filters -\> outputs, so the offsets would be commited after they were consummed, not after they were ingested in elasticsearch.

> [@Chris\_Denneen](#):
>
> Failing to send shouldn't allow it to be "dropped" from retry logic I wouldn't think.

It really depends on why it failed, in this case, logstash didn't failed to send the event, it was correctly sent to elasticsearch and elasticsearch wasn't able to index it and elasticsearch returned an error, if I remind correctly, when you have no write alias elasticsearch returns an error `400`.

When Logstash receives a `400` response for elasticsearch it will per default drop the event and continue on the next one, for `400` and `404` responses you can however enable the [Dead letter queue](https://www.elastic.co/guide/en/logstash/current/dead-letter-queues.html#dead-letter-queues) to store those events and reprocess them using the `dead_letter_queue` input, but this need to be configured upfront.

It seems that this is what happened, elasticsearch returned `400` for logstash requests and the events were dropped until you fixed the write alias and they stated to be indexed.

If your Elasticsearch went down completely for example, Logstash would stop consuming events from Kafka after the queue was full.

> [@Chris\_Denneen](#):
>
> Is there any sort of improvement that can be done on the input that it wouldn't mark it in offset as complete until it gets confirmation about "output" being successful?

For this specific case enabling Dead Letter Queue would help because your events would be stored in the DLQ and you would be able to process them again, but this can lead to other issues, if your DLQ is full it may crash your entire logstash and impact other pipelines.

---

<div class="post-metadata">

**Author:** ![leandrojmp](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/leandrojmp/32/107231_2.png) [@leandrojmp](https://discuss.elastic.co/u/leandrojmp)\
**Post date:** [February 11, 2023, 5:07pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/6 "2023-02-11T17:07:55Z")

</div>

> [@Chris\_Denneen](#):
>
> Do you have any sort of example on the logstash pipeline filter to drop events outside those 2 offsets

As you already saw, you would need to get the offset from the metadata fields, normally I have this on my kafka pipelines:

```auto
filter {
    mutate {
        add_field => {
            "[@metadata][kafka][offset]" => "kafkaOffset"
        }
    }
}

```

So to allow only events betwee two offsets, I think that this would work, but I didn't tested:

```auto
filter {
    mutate {
        convert => {
            "kafkaOffset" => "integer"
        }
    }
    if [kafkaOffset] <= last_offset_indexed or [kafkaOffset] >= first_offset_indexed_after_fix {
        drop {}
    }
}

```

---

<div class="post-metadata">

**Author:** ![Chris\_Denneen](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chris_denneen/32/1254_2.png) [@Chris\_Denneen](https://discuss.elastic.co/u/Chris_Denneen)\
**Post date:** [February 11, 2023, 5:26pm UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/7 "2023-02-11T17:26:32Z")

</div>

DLQ doesn’t support NFS but I’m currently running LS inside Kubernetes so I will probably use EFS to have path for all my index pipelines to write their DLQ. I can then spin up a separate Deployment for processing DLQ that could have been written out from any of the different pipelines I have or N replicas of a particular pipeline. I get NFS is a performance hit but in this world writing to an EBS PV instead of EFS I’d have to unmount from N different pipelines in order to read them. Not sure I see another ReadWriteMany option.

---

<div class="post-metadata">

**Author:** ![Chris\_Denneen](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chris_denneen/32/1254_2.png) [@Chris\_Denneen](https://discuss.elastic.co/u/Chris_Denneen)\
**Post date:** [February 12, 2023, 2:04am UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/8 "2023-02-12T02:04:14Z")

</div>

> [@leandrojmp](#):
>
> > [@Chris\_Denneen](#):
> >
> > Do you have any sort of example on the logstash pipeline filter to drop events outside those 2 offsets
> 
> As you already saw, you would need to get the offset from the metadata fields, normally I have this on my kafka pipelines:
> 
> ```auto
> filter {
> mutate {
> add_field => {
> "[@metadata][kafka][offset]" => "kafkaOffset"
> }
> }
> }
> 
> ```
> 
> So to allow only events betwee two offsets, I think that this would work, but I didn't tested:
> 
> ```auto
> filter {
> mutate {
> convert => {
> "kafkaOffset" => "integer"
> }
> }
> if [kafkaOffset] <= last_offset_indexed or [kafkaOffset] >= first_offset_indexed_after_fix {
> drop {}
> }
> }
> 
> ```

You aren't capturing the partition?

Just grabbed 1 topic now that I've recently added these metadata fields and it seems the offset is partition specific so probably need like an AND to that if condition?

```auto
Feb 11, 2023 @ 21:05:55.604	745555988	filebeat-prod-apjson-apmediaapi	0
Feb 11, 2023 @ 21:05:55.584	745555896	filebeat-prod-apjson-apmediaapi	0
Feb 11, 2023 @ 21:05:55.577	747331699	filebeat-prod-apjson-apmediaapi	1
Feb 11, 2023 @ 21:05:55.577	747331724	filebeat-prod-apjson-apmediaapi	1
Feb 11, 2023 @ 21:05:55.542	751208055	filebeat-prod-apjson-apmediaapi	2
Feb 11, 2023 @ 21:05:55.486	751208101	filebeat-prod-apjson-apmediaapi	2

```

based on as timestamp is going up the offset is getting less would make me think they are partition specific?

---

<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:** [March 12, 2023, 2:04am UTC](https://discuss.elastic.co/t/kafka-input-resync-missing-items-from-topic/325294/9 "2023-03-12T02:04:38Z")

</div>

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