diff options
Diffstat (limited to 'jb/flow')
| -rw-r--r-- | jb/flow/events.py | 16 |
1 files changed, 11 insertions, 5 deletions
diff --git a/jb/flow/events.py b/jb/flow/events.py index f224c01..bee63fd 100644 --- a/jb/flow/events.py +++ b/jb/flow/events.py @@ -10,7 +10,7 @@ from jb.config import ( CONSUMER_NAME, JB_EVENTS_STREAM, ) -from jb.decorators import LOG, REDIS +from jb.decorators import LOG, get_redis from jb.flow.assignment_tasks import process_assignment_submitted from jb.flow.monitoring import emit_error_event from jb.models.event import MTurkEvent @@ -39,7 +39,10 @@ def process_mturk_events(executor: Executor): def create_consumer_group(): try: - REDIS.xgroup_create(JB_EVENTS_STREAM, CONSUMER_GROUP, id="0", mkstream=True) + redis_client = get_redis() + redis_client.xgroup_create( + JB_EVENTS_STREAM, CONSUMER_GROUP, id="0", mkstream=True + ) except redis.exceptions.ResponseError as e: if "BUSYGROUP Consumer Group name already exists" in str(e): pass # group already exists @@ -48,7 +51,8 @@ def create_consumer_group(): def process_mturk_events_chunk(executor: Executor) -> int | None: - msgs_raw = REDIS.xreadgroup( + redis_client = get_redis() + msgs_raw = redis_client.xreadgroup( groupname=CONSUMER_GROUP, consumername=CONSUMER_NAME, streams={JB_EVENTS_STREAM: ">"}, @@ -71,7 +75,7 @@ def process_mturk_events_chunk(executor: Executor) -> int | None: ) else: LOG.info(f"Discarding {event}") - REDIS.xdel(JB_EVENTS_STREAM, msg_id) + redis_client.xdel(JB_EVENTS_STREAM, msg_id) futures.wait(fs, timeout=60) return len(msgs) @@ -80,6 +84,8 @@ def process_mturk_events_chunk(executor: Executor) -> int | None: def process_assignment_submitted_event(event: MTurkEvent, msg_id: str): from jb.decorators import AM, AMTM, BM, HM + redis_client = get_redis() + try: process_assignment_submitted(amtm=AMTM, am=AM, hm=HM, bm=BM, event=event) except Exception as e: @@ -89,4 +95,4 @@ def process_assignment_submitted_event(event: MTurkEvent, msg_id: str): amt_hit_type_id=event.amt_hit_type_id, ) - REDIS.xackdel(JB_EVENTS_STREAM, CONSUMER_GROUP, msg_id) + redis_client.xackdel(JB_EVENTS_STREAM, CONSUMER_GROUP, msg_id) |
