Skip to content

Commit 837eea6

Browse files
bafultonclaude
andcommitted
feat: add --generic-hostname-worker-task-metric flag
Mirror --generic-hostname-task-sent-metric for the worker-executed metrics: label the celery_task_* counters and celery_task_runtime with a generic hostname instead of the executing worker's, to bound label cardinality when workers are autoscaled (for example Kubernetes/KEDA pods). track_task_event builds a single labels dict that feeds every counter and the runtime observation, so one gated assignment collapses them together. celery_task_sent keeps its own flag, and celery_worker_up / celery_worker_tasks_active are unaffected since they do not pass through track_task_event. Refs #387 Co-Authored-By: Claude <noreply@anthropic.com>
1 parent ae2a019 commit 837eea6

3 files changed

Lines changed: 99 additions & 1 deletion

File tree

src/cli.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,16 @@ def _eq_sign_separated_argument_to_dict(_ctx, _param, value):
111111
"Knowing which client sent a task might not be useful for many use cases as for example in "
112112
"Kubernetes environments where the client's hostname is a random string.",
113113
)
114+
@click.option(
115+
"--generic-hostname-worker-task-metric",
116+
default=False,
117+
is_flag=True,
118+
help="The celery_task_* counters and celery_task_runtime will be labeled with a generic "
119+
"hostname instead of the executing worker's hostname. This option helps with label "
120+
"cardinality when using a dynamic number of workers, as for example in Kubernetes "
121+
"environments where the worker's hostname is a random string. celery_task_sent is not "
122+
"affected; use --generic-hostname-task-sent-metric for the client side.",
123+
)
114124
@click.option(
115125
"-Q",
116126
"--queues",
@@ -155,6 +165,7 @@ def cli( # pylint: disable=too-many-arguments,too-many-positional-arguments,too
155165
worker_timeout,
156166
purge_offline_worker_metrics,
157167
generic_hostname_task_sent_metric,
168+
generic_hostname_worker_task_metric,
158169
queues,
159170
metric_prefix,
160171
default_queue_name,
@@ -167,6 +178,7 @@ def cli( # pylint: disable=too-many-arguments,too-many-positional-arguments,too
167178
worker_timeout,
168179
purge_offline_worker_metrics,
169180
generic_hostname_task_sent_metric,
181+
generic_hostname_worker_task_metric,
170182
queues,
171183
metric_prefix,
172184
default_queue_name,

src/exporter.py

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ def __init__(
2727
worker_timeout_seconds=5 * 60,
2828
purge_offline_worker_metrics_seconds=10 * 60,
2929
generic_hostname_task_sent_metric=False,
30+
generic_hostname_worker_task_metric=False,
3031
initial_queues=None,
3132
metric_prefix="celery_",
3233
default_queue_name="celery",
@@ -40,6 +41,7 @@ def __init__(
4041
purge_offline_worker_metrics_seconds
4142
)
4243
self.generic_hostname_task_sent_metric = generic_hostname_task_sent_metric
44+
self.generic_hostname_worker_task_metric = generic_hostname_worker_task_metric
4345
self.default_queue_name = default_queue_name
4446

4547
# Static labels
@@ -284,7 +286,10 @@ def track_task_event(self, event):
284286
"queue_name": getattr(task, "queue", self.default_queue_name),
285287
**self.static_label,
286288
}
287-
if event["type"] == "task-sent" and self.generic_hostname_task_sent_metric:
289+
if event["type"] == "task-sent":
290+
if self.generic_hostname_task_sent_metric:
291+
labels["hostname"] = "generic"
292+
elif self.generic_hostname_worker_task_metric:
288293
labels["hostname"] = "generic"
289294

290295
for counter_name, counter in self.state_counters.items():

src/test_metrics.py

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -258,3 +258,84 @@ def succeed():
258258
)
259259
is None
260260
)
261+
262+
263+
def test_worker_generic_task_hostname(threaded_exporter, celery_app, hostname):
264+
threaded_exporter.generic_hostname_worker_task_metric = True
265+
time.sleep(5)
266+
267+
@celery_app.task
268+
def succeed():
269+
pass
270+
271+
succeed.apply_async()
272+
273+
with start_worker(celery_app, without_heartbeat=False):
274+
time.sleep(5)
275+
276+
# The worker-executed counter and the runtime histogram carry the generic
277+
# hostname, not the executing worker's.
278+
assert (
279+
threaded_exporter.registry.get_sample_value(
280+
"celery_task_succeeded_total",
281+
labels={
282+
"hostname": "generic",
283+
"name": "src.test_metrics.succeed",
284+
"queue_name": "celery",
285+
},
286+
)
287+
== 1.0
288+
)
289+
assert (
290+
threaded_exporter.registry.get_sample_value(
291+
"celery_task_runtime_count",
292+
labels={
293+
"hostname": "generic",
294+
"name": "src.test_metrics.succeed",
295+
"queue_name": "celery",
296+
},
297+
)
298+
== 1.0
299+
)
300+
# The runtime histogram is only ever touched by observe() on the collapsed
301+
# worker event, so no series exists under the real worker hostname.
302+
assert (
303+
threaded_exporter.registry.get_sample_value(
304+
"celery_task_runtime_count",
305+
labels={
306+
"hostname": hostname,
307+
"name": "src.test_metrics.succeed",
308+
"queue_name": "celery",
309+
},
310+
)
311+
is None
312+
)
313+
314+
# celery_task_sent is client-side and governed by its own flag, so this flag
315+
# leaves it labeled with the real hostname.
316+
assert (
317+
threaded_exporter.registry.get_sample_value(
318+
"celery_task_sent_total",
319+
labels={
320+
"hostname": hostname,
321+
"name": "src.test_metrics.succeed",
322+
"queue_name": "celery",
323+
},
324+
)
325+
== 1.0
326+
)
327+
328+
# celery_worker_up does not pass through track_task_event, so it keeps the real
329+
# per-worker hostname that KEDA's scaler depends on.
330+
assert (
331+
threaded_exporter.registry.get_sample_value(
332+
"celery_worker_up", labels={"hostname": hostname}
333+
)
334+
== 1.0
335+
)
336+
assert (
337+
threaded_exporter.registry.get_sample_value(
338+
"celery_worker_up", labels={"hostname": "generic"}
339+
)
340+
is None
341+
)

0 commit comments

Comments
 (0)