mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-05-22 22:05:13 +02:00
Processors use shared queues, means there can be more than process on a queue to share load (#265)
This commit is contained in:
parent
c603caa3cc
commit
cd9a208432
2 changed files with 4 additions and 0 deletions
|
|
@ -1,5 +1,6 @@
|
||||||
|
|
||||||
from pulsar.schema import JsonSchema
|
from pulsar.schema import JsonSchema
|
||||||
|
import pulsar
|
||||||
from prometheus_client import Histogram, Info, Counter, Enum
|
from prometheus_client import Histogram, Info, Counter, Enum
|
||||||
import time
|
import time
|
||||||
|
|
||||||
|
|
@ -51,6 +52,7 @@ class Consumer(BaseProcessor):
|
||||||
|
|
||||||
self.consumer = self.client.subscribe(
|
self.consumer = self.client.subscribe(
|
||||||
input_queue, subscriber,
|
input_queue, subscriber,
|
||||||
|
consumer_type=pulsar.ConsumerType.Shared,
|
||||||
schema=JsonSchema(input_schema),
|
schema=JsonSchema(input_schema),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
|
|
||||||
from pulsar.schema import JsonSchema
|
from pulsar.schema import JsonSchema
|
||||||
|
import pulsar
|
||||||
from prometheus_client import Histogram, Info, Counter, Enum
|
from prometheus_client import Histogram, Info, Counter, Enum
|
||||||
import time
|
import time
|
||||||
|
|
||||||
|
|
@ -71,6 +72,7 @@ class ConsumerProducer(BaseProcessor):
|
||||||
|
|
||||||
self.consumer = self.client.subscribe(
|
self.consumer = self.client.subscribe(
|
||||||
input_queue, subscriber,
|
input_queue, subscriber,
|
||||||
|
consumer_type=pulsar.ConsumerType.Shared,
|
||||||
schema=JsonSchema(input_schema),
|
schema=JsonSchema(input_schema),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue