From a068ca7127a5322956b1d71059ab6376799404a3 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 4 Jun 2025 10:47:33 +0100 Subject: [PATCH] Add concurrency to consumer spec --- trustgraph-base/trustgraph/base/consumer_spec.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/trustgraph-base/trustgraph/base/consumer_spec.py b/trustgraph-base/trustgraph/base/consumer_spec.py index 93665476..89581b02 100644 --- a/trustgraph-base/trustgraph/base/consumer_spec.py +++ b/trustgraph-base/trustgraph/base/consumer_spec.py @@ -4,10 +4,11 @@ from . consumer import Consumer from . spec import Spec class ConsumerSpec(Spec): - def __init__(self, name, schema, handler): + def __init__(self, name, schema, handler, concurrency = 1): self.name = name self.schema = schema self.handler = handler + self.concurrency = concurrency def add(self, flow, processor, definition): @@ -24,6 +25,7 @@ class ConsumerSpec(Spec): schema = self.schema, handler = self.handler, metrics = consumer_metrics, + concurrency = self.concurrency ) # Consumer handle gets access to producers and other