celery.contrib.migrate
Инструменты миграции сообщений (брокер <-> брокер).
- classcelery.contrib.migrate.State
-
Состояние выполнения миграции.
- count=0
- filtered=0
- propertystrtotal
- total_apx=0
- exceptioncelery.contrib.migrate.StopFiltering
-
Полупредикат, используемый для сигнала об остановке фильтрации.
- celery.contrib.migrate.migrate_task(producer, body_, message, queues=None)
-
Перенести сообщение одной задачи.
- celery.contrib.migrate.migrate_tasks(source, dest, migrate=<function migrate_task>, app=None, queues=None, **kwargs)
-
Перенести задачи из одного брокера в другой.
- celery.contrib.migrate.move(predicate, connection=None, exchange=None, routing_key=None, source=None, app=None, callback=None, limit=None, transform=None, **kwargs)
-
Найти задачи с помощью фильтрации и переместить их в новую очередь.
- Параметры:
-
-
predicate (Callable) –
Функция фильтрации, определяющая, какие сообщения перемещать. Должна принимать стандартную сигнатуру
(body, message), используемую в обратных вызовах потребителей Kombu. Если предикат указывает, что сообщение нужно переместить, он должен вернуть одно из следующих значений:кортеж
(exchange, routing_key)илиэкземпляр
Queueили-
- любое другое истинное значение означает, что будут использованы указанные аргументы
-
exchangeиrouting_key.
connection (kombu.Connection) – Пользовательское подключение.
source – List[Union[str, kombu.Queue]]: Необязательный список исходных очередей, используемых вместо очередей по умолчанию (очередей из
task_queues). Этот список также может содержать экземплярыQueue.exchange (str, kombu.Exchange) – Биржа назначения по умолчанию.
routing_key (str) – Ключ маршрутизации назначения по умолчанию.
limit (int) – Ограничение количества фильтруемых сообщений.
callback (Callable) – Обратный вызов после перемещения сообщения с сигнатурой
(state, body, message).transform (Callable) – Необязательная функция преобразования возвращаемого значения (назначения) функции фильтрации.
-
Также поддерживаются те же именованные аргументы, что и в
start_filter().Например, операцию
move_task_by_id()можно реализовать следующим образом:def is_wanted_task(body, message): if body['id'] == wanted_id: return Queue('foo', exchange=Exchange('foo'), routing_key='foo') move(is_wanted_task)или с помощью преобразования:
def transform(value): if isinstance(value, str): return Queue(value, Exchange(value), value) return value move(is_wanted_task, transform=transform)Примечание
Предикат также может вернуть кортеж
(exchange, routing_key), чтобы указать назначение, в которое следует переместить задачу, или экземплярQueue. Любое другое истинное значение означает, что задача будет перемещена в биржу/по ключу маршрутизации по умолчанию.
- celery.contrib.migrate.move_by_idmap(map, **kwargs)
-
Переместить задачи, сопоставив их с картой
task_id: queue.Здесь
queue— очередь, в которую нужно переместить задачу.Пример
>>> move_by_idmap({ ... '5bee6e82-f4ac-468e-bd3d-13e8600250bc': Queue('name'), ... 'ada8652d-aef3-466b-abd2-becdaf1b82b3': Queue('name'), ... '3a2b140d-7db1-41ba-ac90-c36a0ef4ab1f': Queue('name')}, ... queues=['hipri'])
- celery.contrib.migrate.move_by_taskmap(map, **kwargs)
-
Переместить задачи, сопоставив их с картой
task_name: queue.queue— очередь, в которую нужно переместить задачу.Пример
>>> move_by_taskmap({ ... 'tasks.add': Queue('name'), ... 'tasks.mul': Queue('name'), ... })
- celery.contrib.migrate.move_direct(predicate, connection=None, exchange=None, routing_key=None, source=None, app=None, callback=None, limit=None, *, transform=<function worker_direct>, **kwargs)
-
Найти задачи с помощью фильтрации и переместить их в новую очередь.
- Параметры:
-
-
predicate (Callable) –
Функция фильтрации, определяющая, какие сообщения перемещать. Должна принимать стандартную сигнатуру
(body, message), используемую в обратных вызовах потребителей Kombu. Если предикат указывает, что сообщение нужно переместить, он должен вернуть одно из следующих значений:кортеж
(exchange, routing_key)илиэкземпляр
Queueили-
- любое другое истинное значение означает, что будут использованы указанные аргументы
-
exchangeиrouting_key.
connection (kombu.Connection) – Пользовательское подключение.
source – List[Union[str, kombu.Queue]]: Необязательный список исходных очередей, используемых вместо очередей по умолчанию (очередей из
task_queues). Этот список также может содержать экземплярыQueue.exchange (str, kombu.Exchange) – Биржа назначения по умолчанию.
routing_key (str) – Ключ маршрутизации назначения по умолчанию.
limit (int) – Ограничение количества фильтруемых сообщений.
callback (Callable) – Обратный вызов после перемещения сообщения с сигнатурой
(state, body, message).transform (Callable) – Необязательная функция преобразования возвращаемого значения (назначения) функции фильтрации.
-
Также поддерживаются те же именованные аргументы, что и в
start_filter().Например, операцию
move_task_by_id()можно реализовать следующим образом:def is_wanted_task(body, message): if body['id'] == wanted_id: return Queue('foo', exchange=Exchange('foo'), routing_key='foo') move(is_wanted_task)или с помощью преобразования:
def transform(value): if isinstance(value, str): return Queue(value, Exchange(value), value) return value move(is_wanted_task, transform=transform)Примечание
Предикат также может вернуть кортеж
(exchange, routing_key), чтобы указать назначение, в которое следует переместить задачу, или экземплярQueue. Любое другое истинное значение означает, что задача будет перемещена в биржу/по ключу маршрутизации по умолчанию.
- celery.contrib.migrate.move_direct_by_id(task_id, dest, **kwargs)
-
Найти задачу по идентификатору и переместить её в другую очередь.
- Параметры:
-
task_id (str) – Идентификатор задачи, которую нужно найти и переместить.
dest – (str, kombu.Queue): Очередь назначения.
transform (Callable) – Необязательная функция преобразования возвращаемого значения (назначения) функции фильтрации.
**kwargs (Any) – Также поддерживаются те же именованные аргументы, что и в
move().
- celery.contrib.migrate.move_task_by_id(task_id, dest, **kwargs)
-
Найти задачу по идентификатору и переместить её в другую очередь.
- Параметры:
-
task_id (str) – Идентификатор задачи, которую нужно найти и переместить.
dest – (str, kombu.Queue): Очередь назначения.
transform (Callable) – Необязательная функция преобразования возвращаемого значения (назначения) функции фильтрации.
**kwargs (Any) – Также поддерживаются те же именованные аргументы, что и в
move().
- celery.contrib.migrate.republish(producer, message, exchange=None, routing_key=None, remove_props=None)
-
Повторно опубликовать сообщение.
- celery.contrib.migrate.start_filter(app, conn, filter, limit=None, timeout=1.0, ack_messages=False, tasks=None, queues=None, callback=None, forever=False, on_declare_queue=None, consume_from=None, state=None, accept=None, **kwargs)
-
Фильтровать задачи.
- celery.contrib.migrate.task_id_eq(task_id, body, message)
-
Вернуть true, если идентификатор задачи равен task_id’.
- celery.contrib.migrate.task_id_in(ids, body, message)
-
Вернуть true, если идентификатор задачи входит в набор ids’.
Copyright © 2017-2026 Asif Saif Uddin, core team & contributors. All rights reserved.
Celery is licensed under The BSD License (3 Clause, also known as the new BSD license). The license is an OSI approved Open Source license and is GPL-compatible.
https://docs.celeryq.dev/en/stable/reference/celery.contrib.migrate.html