# Indexing Progressively Slows on Thrift Input Stream

**URL:** <https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141>\
**Category:** Elasticsearch\
**Created:** [March 22, 2011, 2:39am UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141 "2011-03-22T02:39:03Z")\
**Posts on this page:** 5\
**Page:** 1

<div class="post-metadata">

**Author:** ![Wolf\_2](https://avatars.discourse-cdn.com/v4/letter/w/ac8455/32.png) [@Wolf\_2](https://discuss.elastic.co/u/Wolf_2)\
**Post date:** [March 22, 2011, 2:39am UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141/1 "2011-03-22T02:39:03Z")

</div>

Hi,

I currently have ES installed on a Xen cluster of 14 instances:

400 GB Disk space  
2 Intel cores  
4 GB RAM

All instances are on the same network segment.

I am storing documents with mapping identical to twitter with a  
collator and analyzer operating on the text field.

I am running 14 python threads each writing to one of the instances on  
port 9500 using the thrift transport.

For each thrift transport open/close transaction I am passing 350  
documents.

For 56 shards and no replication on the index.

When I begin the build of the index I am writing ~4500 documents/sec  
but this performance degrades hence I am writing 5M documents at an  
average rate of 1400 documents/sec and 10M at \< 1000 documents per  
second.

When I double the number of shards I get errors in writing.

Any suggestions?

Regards,

Wolf

index setup:

curl -XPUT '[http://192.168.0.129:9200/twitter](http://192.168.0.129:9200/twitter)' -d  
'{ "index" : {  
"numberOfShards" : 42,  
"numberOfReplicas" : 0,  
"analysis" : {  
"analyzer" : {  
"collation" : {  
"tokenizer" : "keyword",  
"filter" : ["myCollator"]  
},  
"my\_analyzer" : {  
"type" : "igo"  
}  
},  
"filter" : {  
"myCollator" : {  
"type" : "icu\_collation",  
"language" : "ja"  
}  
}  
}  
}  
}'

curl -XPUT '[http://192.168.0.129:9200/twitter/\_mapping](http://192.168.0.129:9200/twitter/_mapping)' -d  
@tweet\_mapping

python class for writing tweet docs

class DocWriter:  
def **init** (self, host, data\_file, initial\_index):  
self.index = initial\_index  
self.uri = "/twitter/tweet/"  
self.host = "192.168.0."+str(host)  
self.data\_file = data\_file  
self.socket = TSocket.TSocket("192.168.0.190", 9500)  
self.transport = TTransport.TBufferedTransport(self.socket)  
self.protocol = TBinaryProtocolAccelerated(self.transport)  
self.client = Rest.Client(self.protocol)  
def write\_file\_to\_es(self):  
with open(self.data\_file, 'r') as f:  
self.transport.open()  
line = f.readline()  
while 0 \< len(line):  
self.index = self.index + 1  
line = json.loads(line.rstrip("\n"), object\_hook=self.fix\_obj)  
request = RestRequest(method=1, uri=self.uri + str(self.index),  
headers={}, body=json.dumps(line))  
response = self.client.execute(request)  
response = json.loads(response.body)  
try: response['ok']  
except NameError:  
print response  
self.transport.close()  
f.close  
return self.index  
line = f.readline()  
f.closed  
self.transport.close()  
return self.index  
def fix\_obj(self, obj):  
if 'retweeted\_status' in obj:  
if 'retweet\_count' in obj['retweeted\_status']:

# print(obj['retweeted\_status']['retweet\_count'])

```
obj['retweeted_status']['retweet_count'] =

```

str(obj['retweeted\_status']['retweet\_count'])  
obj['retweeted\_status']=obj['retweeted\_status']  
elif 'retweet\_count' in obj:  
obj['retweet\_count'] = str(obj['retweet\_count'])  
return obj

---

<div class="post-metadata">

**Author:** ![kimchy](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/kimchy/32/44952_2.png) [@kimchy](https://discuss.elastic.co/u/kimchy)\
**Post date:** [March 22, 2011, 8:24am UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141/2 "2011-03-22T08:24:27Z")

</div>

Few things that might affect that:

1. Have you made sure you increased the file descriptor limit?
2. If you open and close a connection every 350 docs, from 15 processes, all working against a single server, you create many sockets. When you close a socket, it does not really gets closed, it ends up in TIME\_WAIT state. And then, if you load the system enough, the OS will throttle opening sockets until others have been recycled. I suggest you use the same connection through the indexing process, and, work against several hosts.
3. It might be a problem in thrift, you might want to consider using HTTP just to check. Use pyes if you use python, its a good client for elasticsearch.

Also, did you see exceptions in elasticsearch logs, or in your responses from it?

We can also talk about options of how to improve perf when it comes to elasticsearch, but I am not sure you hit it yet. Using the bulk API will help the most in your case.  
On Tuesday, March 22, 2011 at 4:39 AM, Wolf wrote:  
Hi,

> I currently have ES installed on a Xen cluster of 14 instances:
> 
> 400 GB Disk space  
> 2 Intel cores  
> 4 GB RAM
> 
> All instances are on the same network segment.
> 
> I am storing documents with mapping identical to twitter with a  
> collator and analyzer operating on the text field.
> 
> I am running 14 python threads each writing to one of the instances on  
> port 9500 using the thrift transport.
> 
> For each thrift transport open/close transaction I am passing 350  
> documents.
> 
> For 56 shards and no replication on the index.
> 
> When I begin the build of the index I am writing ~4500 documents/sec  
> but this performance degrades hence I am writing 5M documents at an  
> average rate of 1400 documents/sec and 10M at \< 1000 documents per  
> second.
> 
> When I double the number of shards I get errors in writing.
> 
> Any suggestions?
> 
> Regards,
> 
> Wolf
> 
> index setup:
> 
> curl -XPUT '[http://192.168.0.129:9200/twitter](http://192.168.0.129:9200/twitter)' -d  
> '{ "index" : {  
> "numberOfShards" : 42,  
> "numberOfReplicas" : 0,  
> "analysis" : {  
> "analyzer" : {  
> "collation" : {  
> "tokenizer" : "keyword",  
> "filter" : ["myCollator"]  
> },  
> "my\_analyzer" : {  
> "type" : "igo"  
> }  
> },  
> "filter" : {  
> "myCollator" : {  
> "type" : "icu\_collation",  
> "language" : "ja"  
> }  
> }  
> }  
> }  
> }'
> 
> curl -XPUT '[http://192.168.0.129:9200/twitter/\_mapping](http://192.168.0.129:9200/twitter/_mapping)' -d  
> @tweet\_mapping
> 
> python class for writing tweet docs
> 
> class DocWriter:  
> def **init** (self, host, data\_file, initial\_index):  
> self.index = initial\_index  
> self.uri = "/twitter/tweet/"  
> self.host = "192.168.0."+str(host)  
> self.data\_file = data\_file  
> self.socket = TSocket.TSocket("192.168.0.190", 9500)  
> self.transport = TTransport.TBufferedTransport(self.socket)  
> self.protocol = TBinaryProtocolAccelerated(self.transport)  
> self.client = Rest.Client(self.protocol)  
> def write\_file\_to\_es(self):  
> with open(self.data\_file, 'r') as f:  
> self.transport.open()  
> line = f.readline()  
> while 0 \< len(line):  
> self.index = self.index + 1  
> line = json.loads(line.rstrip("\n"), object\_hook=self.fix\_obj)  
> request = RestRequest(method=1, uri=self.uri + str(self.index),  
> headers={}, body=json.dumps(line))  
> response = self.client.execute(request)  
> response = json.loads(response.body)  
> try: response['ok']  
> except NameError:  
> print response  
> self.transport.close()  
> f.close  
> return self.index  
> line = f.readline()  
> f.closed  
> self.transport.close()  
> return self.index  
> def fix\_obj(self, obj):  
> if 'retweeted\_status' in obj:  
> if 'retweet\_count' in obj['retweeted\_status']:
> 
> # print(obj['retweeted\_status']['retweet\_count'])
> 
> obj['retweeted\_status']['retweet\_count'] =  
> str(obj['retweeted\_status']['retweet\_count'])  
> obj['retweeted\_status']=obj['retweeted\_status']  
> elif 'retweet\_count' in obj:  
> obj['retweet\_count'] = str(obj['retweet\_count'])  
> return obj

---

<div class="post-metadata">

**Author:** ![Alberto\_Paro\_2](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/alberto_paro_2/32/1137_2.png) [@Alberto\_Paro\_2](https://discuss.elastic.co/u/Alberto_Paro_2)\
**Post date:** [March 22, 2011, 9:43pm UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141/3 "2011-03-22T21:43:51Z")

</div>

Some tips that can help you:

- using pyes and having thrift you have automatic balancing on es server and a very stable thrift interface threadsafe (pyes share the thrift implementation with the pycassa one, with some tricks for es)

- reading files with readline is very poor performance. Use a bufferreader that wrap the readline.

- do a bulk indexing setting creating record PUT improve performances. The bulk size must be tested to have the "best size": try a bulk inserts with 200, 300, 400, 500 depends of the size of record

- file descriptor and socket problems are very nasty as Shay said. I saw that for every connection netty use a file description, so increment the limit. Also in linux you may set the TIME\_WAIT to 0s, on windows the mininum value is 4s (search on google the regedit key to change)

- reuse the connection. Pyes tries to do this automatically.

- in python socket operations are Gil-Locked, to improve performance use eventlet

If you need some other helps, I'm on IRC [irc.freenode.net](http://irc.freenode.net) #elasticsearch.

Hi,  
Alberto

Il giorno 22/mar/2011, alle ore 09.24, Shay Banon ha scritto:

> Few things that might affect that:
> 
> 1. Have you made sure you increased the file descriptor limit?
> 2. If you open and close a connection every 350 docs, from 15 processes, all working against a single server, you create many sockets. When you close a socket, it does not really gets closed, it ends up in TIME\_WAIT state. And then, if you load the system enough, the OS will throttle opening sockets until others have been recycled. I suggest you use the same connection through the indexing process, and, work against several hosts.
> 3. It might be a problem in thrift, you might want to consider using HTTP just to check. Use pyes if you use python, its a good client for elasticsearch.
> 
> Also, did you see exceptions in elasticsearch logs, or in your responses from it?
> 
> We can also talk about options of how to improve perf when it comes to elasticsearch, but I am not sure you hit it yet. Using the bulk API will help the most in your case.  
> On Tuesday, March 22, 2011 at 4:39 AM, Wolf wrote:
> 
> > Hi,
> > 
> > I currently have ES installed on a Xen cluster of 14 instances:
> > 
> > 400 GB Disk space  
> > 2 Intel cores  
> > 4 GB RAM
> > 
> > All instances are on the same network segment.
> > 
> > I am storing documents with mapping identical to twitter with a  
> > collator and analyzer operating on the text field.
> > 
> > I am running 14 python threads each writing to one of the instances on  
> > port 9500 using the thrift transport.
> > 
> > For each thrift transport open/close transaction I am passing 350  
> > documents.
> > 
> > For 56 shards and no replication on the index.
> > 
> > When I begin the build of the index I am writing ~4500 documents/sec  
> > but this performance degrades hence I am writing 5M documents at an  
> > average rate of 1400 documents/sec and 10M at \< 1000 documents per  
> > second.
> > 
> > When I double the number of shards I get errors in writing.
> > 
> > Any suggestions?
> > 
> > Regards,
> > 
> > Wolf
> > 
> > index setup:
> > 
> > curl -XPUT '[http://192.168.0.129:9200/twitter](http://192.168.0.129:9200/twitter)' -d  
> > '{ "index" : {  
> > "numberOfShards" : 42,  
> > "numberOfReplicas" : 0,  
> > "analysis" : {  
> > "analyzer" : {  
> > "collation" : {  
> > "tokenizer" : "keyword",  
> > ""filter" : ["myCollator"]  
> > }},  
> > "my\_analyzer" : {  
> > ""type" : "igo"  
> > }  
> > },  
> > ""filter" : {  
> > "myCollator" : {  
> > "type" : "icu\_collation",  
> > ""language" : "ja"  
> > }  
> > }  
> > }  
> > }  
> > }'
> > 
> > curl -XPUT '[http://192.168.0.129:9200/twitter/\_mapping](http://192.168.0.129:9200/twitter/_mapping)' -d  
> > @tweet\_mapping
> > 
> > python class for writing tweet docs
> > 
> > class DocWriter:  
> > def **init** (self, host, data\_file, initial\_index):  
> > self.index = initial\_index  
> > self.uri = "/twitter/tweet/"  
> > self.host = "192.168.0."+str(host)  
> > self.data\_file = data\_file  
> > self.socket = TSocket.TSocket("192.168.0.190", 9500)  
> > self.transport = TTransport.TBufferedTransport(self.socket)  
> > self.protocol = TBinaryProtocolAccelerated(self.transport)  
> > self.client = Rest.Client(self.protocol)  
> > def write\_file\_to\_es(self):  
> > with open(self.data\_file, 'r') as f:  
> > self.transport.open()  
> > line = f.readline()  
> > while 0 \< len(line):  
> > self.index = self.index + 1  
> > line = json.loads(line.rstrip("\n"), object\_hook=self.fix\_obj)  
> > request ==3D RestRequest(method=1, uri=self.uri + str(self.index),  
> > headers={}, body=json.dumps(line))  
> > response = self.client.execute(request)  
> > response = json.loads(response.body)  
> > try: response[['ok']  
> > except NameError:  
> > print response  
> > self.transport.close()  
> > f.close  
> > return self.index  
> > line = f.readline()  
> > f.closed  
> > self.transport.close()  
> > return self.index  
> > def fix\_obj(self, obj):  
> > if 'retweeted\_status' in obj:  
> > if 'retweet\_count' in obj['retweeted\_status']:
> > 
> > # print(obj['retweeted\_status']['retweet\_count'])
> > 
> > obj['retweeted\_status']['retweet\_count'] =  
> > str(obj['retweeted\_status']['retweet\_count'])  
> > obj['retweeted\_status']=obj['retweeted\_status']  
> > elif 'retweet\_count' in obj:  
> > obj['retweet\_count'] = str(obj['retweet\_count'])  
> > return obj

---

<div class="post-metadata">

**Author:** ![tsuna](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/tsuna/32/2897_2.png) [@tsuna](https://discuss.elastic.co/u/tsuna)\
**Post date:** [March 27, 2011, 7:23am UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141/4 "2011-03-27T07:23:45Z")

</div>

On Tue, Mar 22, 2011 at 2:43 PM, Alberto Paro [alberto.paro@gmail.com](mailto:alberto.paro@gmail.com) wrote:

> - file descriptor and socket problems are very nasty as Shay said. I saw  
> that for every connection netty use a file description, so increment the

This isn't about Netty, it's just how BSD sockets work. One  
connection = one file descriptor.

> limit. Also in linux you may set the TIME\_WAIT to 0s

Don't do this. There's a reason why TIME\_WAIT exists. Generally  
don't tune any TCP knob unless you really know what you're doing.  
Note: if you haven't read the code of the TCP implementation of the  
Linux kernel, chances are high you don't know what you're doing.  
Don't listen to blog posts and others that recommend tuning up or down  
all the TCP parameters to extreme values, they generally will make  
your problem worse, even if you don't realize it immediately.

This is the type of process you should go through before making any  
change to TCP knobs on Linux:

> **[The "Out of socket memory" error](https://blog.tsunanet.net/2011/03/out-of-socket-memory.html)**
>
> I recently did some work on some of our frontend machines (on which we run Varnish ) at StumbleUpon and decided to track down some of the er...

--  
Benoit "tsuna" Sigoure  
Software Engineer @ [www.StumbleUpon.com](http://www.StumbleUpon.com)

---

<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, 2017, 4:09am UTC](https://discuss.elastic.co/t/indexing-progressively-slows-on-thrift-input-stream/4141/5 "2017-07-06T04:09:11Z")

</div>


