diff --git a/celery_exporter/__main__.py b/celery_exporter/__main__.py index 2737dcd..6803773 100644 --- a/celery_exporter/__main__.py +++ b/celery_exporter/__main__.py @@ -12,7 +12,6 @@ LOG_FORMAT = "[%(asctime)s] %(name)s:%(levelname)s: %(message)s" - @click.command(context_settings={"auto_envvar_prefix": "CELERY_EXPORTER"}) @click.option( "--broker-url", @@ -56,6 +55,12 @@ allow_from_autoenv=False, help="JSON object with additional options passed to the underlying transport.", ) +@click.option( + "--route-options", + type=str, + allow_from_autoenv=False, + help="JSON object with additional options passed to the underlying route.", +) @click.option( "--enable-events", is_flag=True, @@ -108,6 +113,7 @@ def main( max_tasks, namespace, transport_options, + route_options, enable_events, use_ssl, ssl_verify, @@ -137,6 +143,17 @@ def main( ) ) sys.exit(1) + + if route_options: + try: + route_option = json.loads(route_options) + except ValueError: + logging.error( + "Error parsing broker transport options from JSON '{}'".format( + route_options + ) + ) + sys.exit(1) broker_use_ssl = generate_broker_use_ssl( use_ssl, @@ -153,6 +170,7 @@ def main( max_tasks, namespace, transport_options, + route_options, enable_events, broker_use_ssl, ) diff --git a/celery_exporter/core.py b/celery_exporter/core.py index f0612ee..64efdd9 100644 --- a/celery_exporter/core.py +++ b/celery_exporter/core.py @@ -23,6 +23,7 @@ def __init__( transport_options=None, enable_events=False, broker_use_ssl=None, + route_options=None ): self._listen_address = listen_address self._max_tasks = max_tasks @@ -31,6 +32,8 @@ def __init__( self._app = celery.Celery(broker=broker_url, broker_use_ssl=broker_use_ssl) self._app.conf.broker_transport_options = transport_options or {} + if route_options: + self._app.conf.task_routes = route_options def start(self): diff --git a/celery_exporter/utils.py b/celery_exporter/utils.py index 7c2317c..71182d2 100644 --- a/celery_exporter/utils.py +++ b/celery_exporter/utils.py @@ -36,8 +36,11 @@ def get_config(app): routes = conf["task_routes"] res[task_name] = default for i in task_wildcard_names: - if i in routes and "queue" in routes[i]: - res[task_name] = routes[i]["queue"] + if i in routes: + if 'queue' in routes: + res[task_name] = routes[i]['queue'] + else: + res[task_name] = routes[i] break else: res[task_name] = default