A consumer reads records from a stream that never ends (a queue, a socket, a tailing log) and writes them to a warehouse in batches.
Write batched_clean(records, size). records is any iterable of dicts — possibly endless. Return an iterator of batches (lists) of cleaned records:
- •skip a record whose
"user_id" is missing or None; - •skip a record whose
"amount" can't be turned into a number with float(...) (missing, None, or a string like "n/a"); - •a kept record becomes
{"user_id": str(user_id).strip(), "amount": float(amount)}; - •batches hold
size cleaned records, in input order; the last batch may be shorter, and there is never an empty batch; - •
size below 1 raises ValueError (when the function is called or when the first batch is requested — either is fine).
It must work on endless input and read only as far as it has to: producing a batch may pull records from records only until that batch is full.
Python 3.13 in your browser — the standard library plus pandas and numpy; no pip installs.