aboutsummaryrefslogtreecommitdiff
path: root/jb/flow
diff options
context:
space:
mode:
Diffstat (limited to 'jb/flow')
-rw-r--r--jb/flow/events.py16
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)