Extract the processing of individual bulk request to a separate method
This commit is contained in:
@@ -40,6 +40,10 @@ def expand_action(data):
|
|||||||
return action, data.get('_source', data)
|
return action, data.get('_source', data)
|
||||||
|
|
||||||
def _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
def _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
||||||
|
"""
|
||||||
|
Split actions into chunks by number or size, serialize them into strings in
|
||||||
|
the process.
|
||||||
|
"""
|
||||||
bulk_actions = []
|
bulk_actions = []
|
||||||
size, action_count = 0, 0
|
size, action_count = 0, 0
|
||||||
for action, data in actions:
|
for action, data in actions:
|
||||||
@@ -65,6 +69,64 @@ def _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
|||||||
if bulk_actions:
|
if bulk_actions:
|
||||||
yield bulk_actions
|
yield bulk_actions
|
||||||
|
|
||||||
|
def _process_bulk_chunk(client, bulk_actions, raise_on_exception=True, raise_on_error=True, **kwargs):
|
||||||
|
"""
|
||||||
|
Send a bulk request to elasticsearch and process the output.
|
||||||
|
"""
|
||||||
|
# if raise on error is set, we need to collect errors per chunk before raising them
|
||||||
|
errors = []
|
||||||
|
|
||||||
|
try:
|
||||||
|
# send the actual request
|
||||||
|
resp = client.bulk('\n'.join(bulk_actions) + '\n', **kwargs)
|
||||||
|
except TransportError as e:
|
||||||
|
# default behavior - just propagate exception
|
||||||
|
if raise_on_exception:
|
||||||
|
raise e
|
||||||
|
|
||||||
|
# if we are not propagating, mark all actions in current chunk as failed
|
||||||
|
err_message = str(e)
|
||||||
|
exc_errors = []
|
||||||
|
|
||||||
|
# deserialize the data back, thisis expensive but only run on
|
||||||
|
# errors if raise_on_exception is false, so shouldn't be a real
|
||||||
|
# issue
|
||||||
|
bulk_data = iter(map(client.transport.serializer.loads, bulk_actions))
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
# collect all the information about failed actions
|
||||||
|
action = next(bulk_data)
|
||||||
|
op_type, action = action.popitem()
|
||||||
|
info = {"error": err_message, "status": e.status_code, "exception": e}
|
||||||
|
if op_type != 'delete':
|
||||||
|
info['data'] = next(bulk_data)
|
||||||
|
info.update(action)
|
||||||
|
exc_errors.append({op_type: info})
|
||||||
|
except StopIteration:
|
||||||
|
break
|
||||||
|
|
||||||
|
# emulate standard behavior for failed actions
|
||||||
|
if raise_on_error:
|
||||||
|
raise BulkIndexError('%i document(s) failed to index.' % len(exc_errors), exc_errors)
|
||||||
|
else:
|
||||||
|
for err in exc_errors:
|
||||||
|
yield False, err
|
||||||
|
return
|
||||||
|
|
||||||
|
# go through request-reponse pairs and detect failures
|
||||||
|
for op_type, item in map(methodcaller('popitem'), resp['items']):
|
||||||
|
ok = 200 <= item.get('status', 500) < 300
|
||||||
|
if not ok and raise_on_error:
|
||||||
|
errors.append({op_type: item})
|
||||||
|
|
||||||
|
if ok or not errors:
|
||||||
|
# if we are not just recording all errors to be able to raise
|
||||||
|
# them all at once, yield items individually
|
||||||
|
yield ok, {op_type: item}
|
||||||
|
|
||||||
|
if errors:
|
||||||
|
raise BulkIndexError('%i document(s) failed to index.' % len(errors), errors)
|
||||||
|
|
||||||
def streaming_bulk(client, actions, chunk_size=500, max_chunk_bytes=100 * 1014 * 1024,
|
def streaming_bulk(client, actions, chunk_size=500, max_chunk_bytes=100 * 1014 * 1024,
|
||||||
raise_on_error=True, expand_action_callback=expand_action,
|
raise_on_error=True, expand_action_callback=expand_action,
|
||||||
raise_on_exception=True, **kwargs):
|
raise_on_exception=True, **kwargs):
|
||||||
@@ -122,63 +184,11 @@ def streaming_bulk(client, actions, chunk_size=500, max_chunk_bytes=100 * 1014 *
|
|||||||
should return a tuple containing the action line and the data line
|
should return a tuple containing the action line and the data line
|
||||||
(`None` if data line should be omitted).
|
(`None` if data line should be omitted).
|
||||||
"""
|
"""
|
||||||
serializer = client.transport.serializer
|
|
||||||
actions = map(expand_action_callback, actions)
|
actions = map(expand_action_callback, actions)
|
||||||
|
|
||||||
# if raise on error is set, we need to collect errors per chunk before raising them
|
for bulk_actions in _chunk_actions(actions, chunk_size, max_chunk_bytes, client.transport.serializer):
|
||||||
errors = []
|
for result in _process_bulk_chunk(client, bulk_actions, raise_on_exception, raise_on_error, **kwargs):
|
||||||
|
yield result
|
||||||
for bulk_actions in _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
|
||||||
try:
|
|
||||||
# send the actual request
|
|
||||||
resp = client.bulk('\n'.join(bulk_actions) + '\n', **kwargs)
|
|
||||||
except TransportError as e:
|
|
||||||
# default behavior - just propagate exception
|
|
||||||
if raise_on_exception:
|
|
||||||
raise e
|
|
||||||
|
|
||||||
# if we are not propagating, mark all actions in current chunk as failed
|
|
||||||
err_message = str(e)
|
|
||||||
exc_errors = []
|
|
||||||
|
|
||||||
# deserialize the data back, thisis expensive but only run on
|
|
||||||
# errors if raise_on_exception is false, so shouldn't be a real
|
|
||||||
# issue
|
|
||||||
bulk_data = iter(map(serializer.loads, bulk_actions))
|
|
||||||
while True:
|
|
||||||
try:
|
|
||||||
# collect all the information about failed actions
|
|
||||||
action = next(bulk_data)
|
|
||||||
op_type, action = action.popitem()
|
|
||||||
info = {"error": err_message, "status": e.status_code, "exception": e}
|
|
||||||
if op_type != 'delete':
|
|
||||||
info['data'] = next(bulk_data)
|
|
||||||
info.update(action)
|
|
||||||
exc_errors.append({op_type: info})
|
|
||||||
except StopIteration:
|
|
||||||
break
|
|
||||||
|
|
||||||
# emulate standard behavior for failed actions
|
|
||||||
if raise_on_error:
|
|
||||||
raise BulkIndexError('%i document(s) failed to index.' % len(exc_errors), exc_errors)
|
|
||||||
else:
|
|
||||||
for err in exc_errors:
|
|
||||||
yield False, err
|
|
||||||
continue
|
|
||||||
|
|
||||||
# go through request-reponse pairs and detect failures
|
|
||||||
for op_type, item in map(methodcaller('popitem'), resp['items']):
|
|
||||||
ok = 200 <= item.get('status', 500) < 300
|
|
||||||
if not ok and raise_on_error:
|
|
||||||
errors.append({op_type: item})
|
|
||||||
|
|
||||||
if not errors:
|
|
||||||
# if we are not just recording all errors to be able to raise
|
|
||||||
# them all at once, yield items individually
|
|
||||||
yield ok, {op_type: item}
|
|
||||||
|
|
||||||
if errors:
|
|
||||||
raise BulkIndexError('%i document(s) failed to index.' % len(errors), errors)
|
|
||||||
|
|
||||||
def bulk(client, actions, stats_only=False, **kwargs):
|
def bulk(client, actions, stats_only=False, **kwargs):
|
||||||
"""
|
"""
|
||||||
|
|||||||
Reference in New Issue
Block a user