import abc import io import itertools from collections.abc import Callable, Iterator from concurrent.futures import ThreadPoolExecutor, as_completed import pandas as pd from pandas.io.parsers import TextFileReader from app.adapters.ows_dmp import AudienceShare, AudienceSharePlatform from app.types import ProcessedFans from app.utils import smart_iter class AudienceShareStrategy(abc.ABC): platform: AudienceSharePlatform batch_size: int = 16 @abc.abstractmethod def share( self, audience_share: AudienceShare, file_buffer: io.BufferedRandom ) -> Iterator[ProcessedFans]: pass def process_in_parallel( self, *, df: TextFileReader, process_fn: Callable[[pd.DataFrame, int, bool], ProcessedFans], ) -> Iterator[ProcessedFans]: """Process DataFrame in parallel and yield ProcessedFans.""" df_iter = smart_iter(df, start=1) while True: batch_items = list(itertools.islice(df_iter, self.batch_size)) if not batch_items: break results = {} # Run batch uploads in parallel with ThreadPoolExecutor(max_workers=len(batch_items)) as executor: future_to_batch = { executor.submit(process_fn, data, batch_seq, is_last): batch_seq for batch_seq, is_last, data in batch_items } for future in as_completed(future_to_batch): batch_seq = future_to_batch[future] results[batch_seq] = future.result() processed_fans = 0 last_batch_seq = max(results.keys()) for batch_seq in sorted(results): result = results.pop(batch_seq) processed_fans += result.processed_fans if batch_seq == last_batch_seq: yield ProcessedFans( processed_fans=processed_fans, is_last=result.is_last )