I would like to share my experience when using itertools.tee to split a large size plat text file into multiple csv files from/to s3 at Python 3.6.9 and 3.7.4 environment.
My data flow are from s3 zipfile, s3fs read iter, map iter for dataclass transform, tee iter, map iter for dataclass filter, loop over the iter and capture data and write to s3 in csv format with s3fs write and/or local write and s3fs put to s3.
The itertools.tee was failed on the zipfile process stack.
Above, Dror Speiser, safetee worked fine, but memory usage increased for any unbalances between tee object as dataset not good distribution or processing delays.
Also, it was not properly work with multiprocessing-logging, might be related this bug: https://bugs.python.org/issue34410
Below code is just to add simple flow control in between tee object to prevent memory increment and OOM Killer situation.
Hope to be helpful for future reference.
import time
import threading
import logging
from itertools import tee
from collections import Counter
logger = logging.getLogger(__name__)
FLOW_WAIT_GAP = 1000 # flow gap for waiting
FLOW_WAIT_TIMEOUT = 60.0 # flow wait timeout
class Safetee:
"""tee object wrapped to make it thread-safe and flow controlled"""
def __init__(self, teeobj, lock, flows, teeidx):
self.teeobj = teeobj
self.lock = lock
self.flows = flows
self.mykey = teeidx
self.logcnt = 0
def __iter__(self):
return self
def __next__(self):
waitsec = 0.0
while True:
with self.lock:
flowgap = self.flows[self.mykey] - self.flows[len(self.flows) - 1]
if flowgap < FLOW_WAIT_GAP or waitsec > FLOW_WAIT_TIMEOUT:
nextdata = next(self.teeobj)
self.flows[self.mykey] += 1
return nextdata
waitthis = min(flowgap / FLOW_WAIT_GAP, FLOW_WAIT_TIMEOUT / 3)
waitsec += waitthis
time.sleep(waitthis)
if waitsec > FLOW_WAIT_TIMEOUT and self.logcnt < 5:
self.logcnt += 1
logger.debug(f'tee wait seconds={waitsec:.2f}, mykey={self.mykey}, flows={self.flows}')
def __copy__(self):
return Safetee(self.teeobj.__copy__(), self.lock, self.flows, self.teeidx)
def safetee(iterable, n=2):
"""tuple of n independent thread-safe and flow controlled iterators"""
lock = threading.Lock()
flows = Counter()
return tuple(Safetee(teeobj, lock, flows, teeidx) for teeidx, teeobj in enumerate(tee(iterable, n)))