Add optional callback to ChunkedInsert
This commit adds an optional `callback` argument to the `ChunkedInsert`
object. This callback is a callable object which gets called before the
chunked insert happens. This is useful for clearing any local caches
that may be in place to deal with the eventual consistency resulting
from the delayed nature of the chunked inserts.
For example,
```
cache = set()
chunked_table = ChunkedInsert(table, callback=lambda queue: cache.clear())
while True:
data = get_data_id()
key = data['key']
if key in cache or table.find_one(key=key)
continue
cache.add(key)
chunked_table.insert(data)
```
This commit is contained in:
+14
-2
@@ -1,23 +1,35 @@
|
||||
|
||||
|
||||
class InvalidCallback(ValueError):
|
||||
pass
|
||||
|
||||
|
||||
class ChunkedInsert(object):
|
||||
"""Batch up insert operations
|
||||
with ChunkedStorer(my_table) as storer:
|
||||
table.insert(row)
|
||||
|
||||
Rows will be inserted in groups of 1000
|
||||
Rows will be inserted in groups of `chunksize` (defaulting to 1000). An
|
||||
optional callback can be provided that will be called before the insert.
|
||||
This callback takes one parameter which is the queue which is about to be
|
||||
inserted into the database
|
||||
"""
|
||||
|
||||
def __init__(self, table, chunksize=1000):
|
||||
def __init__(self, table, chunksize=1000, callback=None):
|
||||
self.queue = []
|
||||
self.fields = set()
|
||||
self.table = table
|
||||
self.chunksize = chunksize
|
||||
if callback and not callable(callback):
|
||||
raise InvalidCallback
|
||||
self.callback = callback
|
||||
|
||||
def flush(self):
|
||||
for item in self.queue:
|
||||
for field in self.fields:
|
||||
item[field] = item.get(field)
|
||||
if self.callback is not None:
|
||||
self.callback(self.queue)
|
||||
self.table.insert_many(self.queue)
|
||||
self.queue = []
|
||||
|
||||
|
||||
Reference in New Issue
Block a user