@@ -1,7 +1,6 @@
|
|||||||
from __future__ import unicode_literals
|
from __future__ import unicode_literals
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
from itertools import islice
|
|
||||||
from operator import methodcaller
|
from operator import methodcaller
|
||||||
|
|
||||||
from ..exceptions import ElasticsearchException, TransportError
|
from ..exceptions import ElasticsearchException, TransportError
|
||||||
@@ -40,22 +39,35 @@ def expand_action(data):
|
|||||||
|
|
||||||
return action, data.get('_source', data)
|
return action, data.get('_source', data)
|
||||||
|
|
||||||
def _chunk_actions(actions, chunk_size, serializer):
|
def _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
||||||
while True:
|
|
||||||
bulk_actions = []
|
bulk_actions = []
|
||||||
for action, data in islice(actions, chunk_size):
|
size, action_count = 0, 0
|
||||||
bulk_actions.append(serializer.dumps(action))
|
for action, data in actions:
|
||||||
|
action = serializer.dumps(action)
|
||||||
|
cur_size = len(action) + 1
|
||||||
|
|
||||||
if data is not None:
|
if data is not None:
|
||||||
bulk_actions.append(serializer.dumps(data))
|
data = serializer.dumps(data)
|
||||||
|
cur_size += len(data) + 1
|
||||||
|
|
||||||
if not bulk_actions:
|
# full chunk, send it and start a new one
|
||||||
return
|
if bulk_actions and (size + cur_size > max_chunk_bytes or action_count == chunk_size):
|
||||||
|
yield bulk_actions
|
||||||
|
bulk_actions = []
|
||||||
|
size, action_count = 0, 0
|
||||||
|
|
||||||
|
bulk_actions.append(action)
|
||||||
|
if data is not None:
|
||||||
|
bulk_actions.append(data)
|
||||||
|
size += cur_size
|
||||||
|
action_count += 1
|
||||||
|
|
||||||
|
if bulk_actions:
|
||||||
yield bulk_actions
|
yield bulk_actions
|
||||||
|
|
||||||
def streaming_bulk(client, actions, chunk_size=500, raise_on_error=True,
|
def streaming_bulk(client, actions, chunk_size=500, max_chunk_bytes=100 * 1014 * 1024,
|
||||||
expand_action_callback=expand_action, raise_on_exception=True,
|
raise_on_error=True, expand_action_callback=expand_action,
|
||||||
**kwargs):
|
raise_on_exception=True, **kwargs):
|
||||||
"""
|
"""
|
||||||
Streaming bulk consumes actions from the iterable passed in and yields
|
Streaming bulk consumes actions from the iterable passed in and yields
|
||||||
results per action. For non-streaming usecases use
|
results per action. For non-streaming usecases use
|
||||||
@@ -101,6 +113,7 @@ def streaming_bulk(client, actions, chunk_size=500, raise_on_error=True,
|
|||||||
:arg client: instance of :class:`~elasticsearch.Elasticsearch` to use
|
:arg client: instance of :class:`~elasticsearch.Elasticsearch` to use
|
||||||
:arg actions: iterable containing the actions to be executed
|
:arg actions: iterable containing the actions to be executed
|
||||||
:arg chunk_size: number of docs in one chunk sent to es (default: 500)
|
:arg chunk_size: number of docs in one chunk sent to es (default: 500)
|
||||||
|
:arg max_chunk_bytes: the maximum size of the request in bytes (default: 100MB)
|
||||||
:arg raise_on_error: raise ``BulkIndexError`` containing errors (as `.errors`)
|
:arg raise_on_error: raise ``BulkIndexError`` containing errors (as `.errors`)
|
||||||
from the execution of the last chunk when some occur. By default we raise.
|
from the execution of the last chunk when some occur. By default we raise.
|
||||||
:arg raise_on_exception: if ``False`` then don't propagate exceptions from
|
:arg raise_on_exception: if ``False`` then don't propagate exceptions from
|
||||||
@@ -115,7 +128,7 @@ def streaming_bulk(client, actions, chunk_size=500, raise_on_error=True,
|
|||||||
# if raise on error is set, we need to collect errors per chunk before raising them
|
# if raise on error is set, we need to collect errors per chunk before raising them
|
||||||
errors = []
|
errors = []
|
||||||
|
|
||||||
for bulk_actions in _chunk_actions(actions, chunk_size, serializer):
|
for bulk_actions in _chunk_actions(actions, chunk_size, max_chunk_bytes, serializer):
|
||||||
try:
|
try:
|
||||||
# send the actual request
|
# send the actual request
|
||||||
resp = client.bulk('\n'.join(bulk_actions) + '\n', **kwargs)
|
resp = client.bulk('\n'.join(bulk_actions) + '\n', **kwargs)
|
||||||
|
|||||||
@@ -1,9 +1,21 @@
|
|||||||
from datetime import datetime
|
|
||||||
|
|
||||||
from elasticsearch import helpers, TransportError
|
from elasticsearch import helpers, TransportError
|
||||||
|
from elasticsearch.serializer import JSONSerializer
|
||||||
|
|
||||||
from . import ElasticsearchTestCase
|
from . import ElasticsearchTestCase
|
||||||
from ..test_cases import SkipTest
|
from ..test_cases import SkipTest, TestCase
|
||||||
|
|
||||||
|
|
||||||
|
class TestChunkActions(TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
super(TestChunkActions, self).setUp()
|
||||||
|
self.actions = [({'index': {}}, {'some': 'data', 'i': i}) for i in range(100)]
|
||||||
|
|
||||||
|
def test_chunks_are_chopped_by_byte_size(self):
|
||||||
|
self.assertEquals(100, len(list(helpers._chunk_actions(self.actions, 100000, 1, JSONSerializer()))))
|
||||||
|
|
||||||
|
def test_chunks_are_chopped_by_chunk_size(self):
|
||||||
|
self.assertEquals(10, len(list(helpers._chunk_actions(self.actions, 10, 99999999, JSONSerializer()))))
|
||||||
|
|
||||||
|
|
||||||
class FailingBulkClient(object):
|
class FailingBulkClient(object):
|
||||||
def __init__(self, client, fail_at=1):
|
def __init__(self, client, fail_at=1):
|
||||||
|
|||||||
Reference in New Issue
Block a user