Elasticsearch parallel_bulk helper function in python is throwing error for a partial failure when consuming response

Viewed 253

I'm trying to test partial failure scenarios using elasticsearch parallel_bulk function in python. I have the an index with the following mapping

PUT test
{
  "mappings": {
    "dynamic": "strict",
    "properties": {
      "random": {
        "type": "keyword"
      }
    }
  }
}

I wrote the below code where I try to do bulk insert with one doc having an extra field which is not in the mapping and since the mapping is strict, insert should fail for this one doc

from elasticsearch import Elasticsearch
from elasticsearch.helpers import parallel_bulk

def index_bulk(payload: list, index_name: str):
    for data in payload:
        yield {
            '_op_type': 'index',
            '_index': index_name,
            '_source': data
        }

conn_params = {
    'hosts': [{'host': 'localhost', 'port': 9200}],
}
es = Elasticsearch(**conn_params)
docs = [{ "random": f"value_{i}" } for i in range(3)]
docs.append({
    "random": "value",
    "field": "test"
})

response = parallel_bulk(
    client=es,
    actions=index_bulk(docs, 'test')
)

The above code works and doesn't throw any errors. But now I try to consume the response generator to see if there are any failures using the following code

for success, info in response:
    if not success: print('Doc failed', info)

And instead of printing the failed docs, it's throwing the error below

Traceback (most recent call last):
  File "test_bulk_es.py", line 28, in <module>
    for success, info in response:
  File "C:\Users\username\project\venv\lib\site-packages\elasticsearch\helpers\actions.py", line 365, in parallel_bulk 
    actions, chunk_size, max_chunk_bytes, client.transport.serializer
  File "C:\Users\username\AppData\Local\Programs\Python\Python36\Lib\multiprocessing\pool.py", line 735, in next
    raise value
  File "C:\Users\username\AppData\Local\Programs\Python\Python36\Lib\multiprocessing\pool.py", line 119, in worker
    result = (True, func(*args, **kwds))
  File "C:\Users\username\project\venv\lib\site-packages\elasticsearch\helpers\actions.py", line 361, in <lambda>      
    client, bulk_chunk[1], bulk_chunk[0], *args, **kwargs
  File "C:\Users\username\project\venv\lib\site-packages\elasticsearch\helpers\actions.py", line 158, in _process_bulk_chunk
    raise BulkIndexError("%i document(s) failed to index." % len(errors), errors)
elasticsearch.helpers.errors.BulkIndexError: ('1 document(s) failed to index.', [{'index': {'_index': 'test', '_type': '_doc', '_id': 'km93f30Ba79etmZE_XiS', 'status': 400, 'error': {'type': 'strict_dynamic_mapping_exception', 'reason': 'mapping set to strict, dynamic introduction of [field] within [_doc] is not allowed'}, 'data': {'random': 'value', 'field': 'test'}}}])

I don't want it to throw any error. I just want to know if there any partial failure and I assumed that's how the above code works from what I read in the documentation

enter image description here

Am I doing something wrong?

I'm using elasticsearch==7.1.0

1 Answers

The elasticsearch-py source for parallel_bulk says you can pass raise_on_error=False to prevent this. Then your if not success block will work.

Be aware that the text of the original document is only available in the error when raise_on_error is True, but that may still be better than getting the full text of 200 failed documents in your logs and error emails.

Related