Batch processing
There is no separate batch API: evaluation is fast enough that a loop over evaluate() is the batch API. This page covers how to structure that loop so the output is auditable and the throughput scales.
About the example
This page uses the built-in data-protection pack so you can run everything without writing a pack first. The same calls work for any decision pack: eligibility, limits, routing and so on. See Core concepts.
Evaluate per field, not per record#
Serializing a whole record and evaluating it once tells you that something in the record broke a rule. Evaluating each text field tells you which field, and lets you apply the remedy to that field only. With the example pack, that means redacting only the field that needs it:
import csv
from helixor_runtime import HelixorEngine
engine = HelixorEngine()
TEXT_FIELDS = ["notes", "contact", "comment"]
def process(row: dict) -> tuple[dict, list[dict]]:
out, audit = dict(row), []
for field in TEXT_FIELDS:
value = row.get(field) or ""
if not value:
continue
r = engine.evaluate(value)
if r.invariants_passed:
continue
fatal = any(t.severity == "FATAL" for t in r.triggers)
out[field] = "[WITHHELD]" if fatal or not r.action.startswith("redact_") else r.remedy.clean_text
audit.append({"field": field, "action": r.action, "receipt": r.receipt_hash,
"rules": ";".join(t.rule_id for t in r.triggers)})
return out, audit
with open("in.csv") as src, open("out.csv", "w", newline="") as dst:
reader = csv.DictReader(src)
writer = csv.DictWriter(dst, fieldnames=reader.fieldnames)
writer.writeheader()
for row in reader:
clean, audit = process(row)
writer.writerow(clean)
for entry in audit:
audit_log.write(entry) # your append-only audit sink
A field with a fatal finding is withheld entirely rather than remedied, which matches what an online check would do. Change that policy deliberately if your pipeline needs the remedied text instead.
Measure before you scale#
Run the engine-time benchmark on the hardware you will use (download it and _common.py from Reproduce):
python bench_engine_time.py
In the published runs a single process sustained about 35,000 to 42,000 evaluations per second, at about 25 µs per call at the median, on a loaded machine. Treat those as an illustration, not a specification: latency_us only covers evaluation, and your payload lengths change the result. See Performance.
Scale with processes#
Evaluation is CPU-bound Python, so threads do not add throughput and the engine is not designed to be shared between them. Use one engine per process:
from concurrent.futures import ProcessPoolExecutor
from helixor_runtime import HelixorEngine
_engine = None
def _init():
global _engine
_engine = HelixorEngine() # one per worker process
def decide(text: str) -> tuple[str, str]:
r = _engine.evaluate(text)
return r.action, r.receipt_hash
with ProcessPoolExecutor(max_workers=8, initializer=_init) as pool:
for action, receipt in pool.map(decide, texts, chunksize=512):
...
Use a large chunksize: sending one short string to a worker costs more than evaluating it.
Queues and streams of records#
For a message consumer, evaluate inside the handler and publish the audit entry alongside the cleaned message. Keep the receipt with the message ID so a later question ("why was message 8812 changed?") can be answered from the audit log without the original payload.
For a complete worked example with a per-row audit report, see the audit pipeline tutorial.