Spec-Zone.ru › Celery

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. Если предикат указывает, что сообщение нужно переместить, он должен вернуть одно из следующих значений:

    1. кортеж (exchange, routing_key) или

    2. экземпляр Queue или

    3. любое другое истинное значение означает, что будут использованы указанные аргументы

      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. Если предикат указывает, что сообщение нужно переместить, он должен вернуть одно из следующих значений:

    1. кортеж (exchange, routing_key) или

    2. экземпляр Queue или

    3. любое другое истинное значение означает, что будут использованы указанные аргументы

      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

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API