# Throttle the ES-Hadoop write speed

**URL:** <https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702>\
**Category:** Elasticsearch\
**Tags:** es-hadoop\
**Created:** [August 28, 2020, 2:37am UTC](https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702 "2020-08-28T02:37:57Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![chenchuangc](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chenchuangc/32/40017_2.png) [@chenchuangc](https://discuss.elastic.co/u/chenchuangc)\
**Post date:** [August 28, 2020, 2:37am UTC](https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702/1 "2020-08-28T02:37:57Z")

</div>

## 1. problem background

I am using an elasticsearch cluster to save the result from spark2.3.  
Also this cluster offer an online query .  
My spark task is a daily work that will write 6 million records to ES cluster every day .  
now the process of writing to es will use 10min every day , indexing speed is about 10k per second, but during this 10min, there will be some queries spend more the 1 second , but if i test this query in the other time of the day (not this 10min), the query is very fast (about 10 milliseconds).

So i think that maybe the es cluster is over take , and i want to lower the write speed (i can accept longer time for the writing process) so that the online query will be faster.

the ES-Hadoop maven

```auto

        <dependency> <!-- Spark dependency -->
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.3.0</version>
        </dependency>

       <dependency>
            <groupId>org.elasticsearch</groupId>
            <artifactId>elasticsearch-spark-20_2.11</artifactId>
            <version>7.8.1</version>
        </dependency>

```

## 2. The work i have tried

1. use less executors of spark

```auto
--executor-cores 1 --num-executors 1

```

1. reduce the write batch size and close the refresh

```auto
 SparkConf sparkConf = new SparkConf()
               ...
                .set(ConfigurationOptions.ES_BATCH_SIZE_ENTRIES, "50")
                .set(ConfigurationOptions.ES_BATCH_WRITE_REFRESH, "false")

```

but the speed is still too high for me  
i want to know is there any other way to throttle the write speed for ES-Hadoop

sorry for bother you

---

<div class="post-metadata">

**Author:** ![chenchuangc](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chenchuangc/32/40017_2.png) [@chenchuangc](https://discuss.elastic.co/u/chenchuangc)\
**Post date:** [August 31, 2020, 1:32am UTC](https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702/2 "2020-08-31T01:32:51Z")

</div>

any body can help me ?  
if no way , i have to dispose the ES-Hadoop and write to elasticsearch manually by the elasticsearch java client 😂

---

<div class="post-metadata">

**Author:** ![chenchuangc](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/chenchuangc/32/40017_2.png) [@chenchuangc](https://discuss.elastic.co/u/chenchuangc)\
**Post date:** [September 1, 2020, 2:31am UTC](https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702/3 "2020-09-01T02:31:04Z")

</div>

i rewrite the `org.elasticsearch.spark.rdd.EsRDDWriter.write() `

just like this

```auto

/**
  * Created by chencc on 2020/8/31.
  */
@Slf4j
class MyEsDataFrameWriter (schema: StructType, override val serializedSettings: String)
  extends EsRDDWriter[Row](serializedSettings:String) {

  override protected def valueWriter: Class[_ <: ValueWriter[_]] = classOf[DataFrameValueWriter]
  override protected def bytesConverter: Class[_ <: BytesConverter] = classOf[JdkBytesConverter]
  override protected def fieldExtractor: Class[_ <: FieldExtractor] = classOf[DataFrameFieldExtractor]

  override protected def processData(data: Iterator[Row]): Any = { (data.next, schema) }

  override def write(taskContext: TaskContext, data: Iterator[Row]): Unit = {
    val writer = RestService.createWriter(settings, taskContext.partitionId.toLong, -1, log)

    taskContext.addTaskCompletionListener((TaskContext) => writer.close())

    if (runtimeMetadata) {
      writer.repository.addRuntimeFieldExtractor(metaExtractor)
    }

    val counter= new AtomicInteger(0);
    while (data.hasNext) {
      counter.incrementAndGet();
      writer.repository.writeToIndex(processData(data))
      if(counter.get()>=100){
        Thread.sleep(100);
        counter.set(0)
        log.info("batch is 2000 will sleep 50 milliseconds ")
      }
    }
  }
}

```

---

<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:** [September 29, 2020, 2:31am UTC](https://discuss.elastic.co/t/throttle-the-es-hadoop-write-speed/246702/4 "2020-09-29T02:31:04Z")

</div>

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