From ee4834155792fd9ce0bd2909119cfa76c723c90b Mon Sep 17 00:00:00 2001 From: Kirill Nagaitsev Date: Tue, 11 Jan 2022 11:46:30 -0600 Subject: [PATCH 1/2] add rabbitmq queue ttl as env var --- entrypoint.sh | 2 +- funcx_forwarder/forwarder.py | 19 +++++++++++++++++-- funcx_forwarder/service.py | 9 +++++++++ 3 files changed, 27 insertions(+), 3 deletions(-) diff --git a/entrypoint.sh b/entrypoint.sh index 7a0c826..f18efcc 100755 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -25,4 +25,4 @@ if [[ -z "${ADVERTISED_FORWARDER_ADDRESS}" ]]; then fi -forwarder-service -a $ADVERTISED_FORWARDER_ADDRESS -p 8080 --redishost $REDIS_HOST --redisport $REDIS_PORT --rabbitmquri $RABBITMQ_URI -d --endpoint-base-port ${ENDPOINT_BASE_PORT} +forwarder-service -a $ADVERTISED_FORWARDER_ADDRESS -p 8080 --redishost $REDIS_HOST --redisport $REDIS_PORT --rabbitmquri $RABBITMQ_URI -d --endpoint-base-port ${ENDPOINT_BASE_PORT} --rabbitmq_queue_ttl $RABBITMQ_QUEUE_TTL diff --git a/funcx_forwarder/forwarder.py b/funcx_forwarder/forwarder.py index 0dcafeb..259bc60 100644 --- a/funcx_forwarder/forwarder.py +++ b/funcx_forwarder/forwarder.py @@ -106,7 +106,8 @@ def __init__( response_queue, address: str, redis_address: str, - rabbitmq_conn_params, + rabbitmq_conn_params: pika.URLParameters, + rabbitmq_queue_ttl: int, endpoint_ports=(55001, 55002, 55003), redis_port: int = 6379, logging_level=logging.INFO, @@ -131,6 +132,13 @@ def __init__( redis_address : str full address to connect to redis. Required + rabbitmq_conn_params : pika.URLParameters + URL Params for RabbitMQ connection + + rabbitmq_queue_ttl : int + RabbitMQ queue TTL in seconds + (must match websocket service rabbitmq_queue_ttl) + endpoint_ports : (int, int, int) A triplet of ports: (tasks_port, results_port, commands_port) Default: (55001, 55002, 55003) @@ -159,6 +167,7 @@ def __init__( self.address = address self.redis_url = f"{redis_address}:{redis_port}" self.rabbitmq_conn_params = rabbitmq_conn_params + self.rabbitmq_queue_ttl = rabbitmq_queue_ttl self.tasks_port, self.results_port, self.commands_port = endpoint_ports self.connected_endpoints: t.Dict[str, t.Dict[str, t.Any]] = {} self.kill_event = Event() @@ -177,6 +186,7 @@ def __init__( logger.info(f"Initializing forwarder v{funcx_forwarder.__version__}") logger.info(f"Forwarder running on public address: {self.address}") logger.info(f"REDIS url: {self.redis_url}") + logger.info(f"RabbitMQ Queue TTL: {self.rabbitmq_queue_ttl}") logger.info(f"Log level set to {loglevels[logging_level]}") if not os.path.exists(self.keys_dir) or not os.listdir(self.keys_dir): @@ -633,10 +643,15 @@ def handle_results(self): # and the task result will not be acked if this fails task_group_id = task.task_group_id if task_group_id: + # This argument expects milliseconds, so multiply by 1000 + queue_args = { + "x-expires": int(self.rabbitmq_queue_ttl * 1000), + } + connection = pika.BlockingConnection(self.rabbitmq_conn_params) channel = connection.channel() channel.exchange_declare(exchange="tasks", exchange_type="direct") - channel.queue_declare(queue=task_group_id) + channel.queue_declare(queue=task_group_id, arguments=queue_args) channel.queue_bind(task_group_id, "tasks") # important: the FuncX client must be capable of receiving the same diff --git a/funcx_forwarder/service.py b/funcx_forwarder/service.py index 07122fe..25609e7 100644 --- a/funcx_forwarder/service.py +++ b/funcx_forwarder/service.py @@ -194,6 +194,14 @@ def cli_run(): parser.add_argument( "-d", "--debug", action="store_true", help="Enables debug logging" ) + parser.add_argument( + "--rabbitmq_queue_ttl", + required=True, + help=( + "Set RabbitMQ queue TTL in seconds " + "(must match websocket service queue TTL)" + ), + ) parser.add_argument( "-v", "--version", action="store_true", help="Print version information" ) @@ -237,6 +245,7 @@ def cli_run(): args.address, args.redishost, rabbitmq_conn_params, + int(args.rabbitmq_queue_ttl), endpoint_ports=range(args.endpoint_base_port, args.endpoint_base_port + 3), logging_level=logging_level, redis_port=args.redisport, From 68de30ef597bf57f7e228065e37b57fdd003a56e Mon Sep 17 00:00:00 2001 From: Kirill Nagaitsev Date: Tue, 11 Jan 2022 13:15:49 -0600 Subject: [PATCH 2/2] set default queue ttl if not provided --- entrypoint.sh | 8 +++++++- funcx_forwarder/service.py | 2 +- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/entrypoint.sh b/entrypoint.sh index f18efcc..48f9328 100755 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -16,6 +16,12 @@ if [[ -z "${ENDPOINT_BASE_PORT}" ]]; then ENDPOINT_BASE_PORT=55001 fi +if [[ -z "${RABBITMQ_QUEUE_TTL}" ]]; then + RABBITMQ_QUEUE_TTL_OPT="" +else + RABBITMQ_QUEUE_TTL_OPT="--rabbitmq_queue_ttl $RABBITMQ_QUEUE_TTL" +fi + python3 wait_for_redis.py if [[ -z "${ADVERTISED_FORWARDER_ADDRESS}" ]]; then @@ -25,4 +31,4 @@ if [[ -z "${ADVERTISED_FORWARDER_ADDRESS}" ]]; then fi -forwarder-service -a $ADVERTISED_FORWARDER_ADDRESS -p 8080 --redishost $REDIS_HOST --redisport $REDIS_PORT --rabbitmquri $RABBITMQ_URI -d --endpoint-base-port ${ENDPOINT_BASE_PORT} --rabbitmq_queue_ttl $RABBITMQ_QUEUE_TTL +forwarder-service -a $ADVERTISED_FORWARDER_ADDRESS -p 8080 --redishost $REDIS_HOST --redisport $REDIS_PORT --rabbitmquri $RABBITMQ_URI -d --endpoint-base-port ${ENDPOINT_BASE_PORT} $RABBITMQ_QUEUE_TTL_OPT diff --git a/funcx_forwarder/service.py b/funcx_forwarder/service.py index 25609e7..5f8b7f5 100644 --- a/funcx_forwarder/service.py +++ b/funcx_forwarder/service.py @@ -196,7 +196,7 @@ def cli_run(): ) parser.add_argument( "--rabbitmq_queue_ttl", - required=True, + default=604800, help=( "Set RabbitMQ queue TTL in seconds " "(must match websocket service queue TTL)"