Skip to content

kiq jobs with unique names, so they won't be created twice #271

Description

@gegnew

I can't seem to find anything about this in the docs. I'd like to assign a job (not a task) a unique name, so that if I were to try to kiq a job with the same parameters twice, the broker would ignore the duplicate job. Something like:

@broker.task
async def foo(param):
    print(param)
    await asyncio.sleep(1)

if __name__ == __main__:
    params = ['a', 'a', 'b', 'c']  # the second "a" should be skipped
    for p in params:
        await foo.kiq(p, job_id=f"job_{p}")

>>> "a"
>>> "b"
>>> "c"

Is this possible?

Activity

  1. gegnew commented on Jan 5, 2024

    @gegnew
    Author

    Clearly the AsyncBroker can take an alternative task_id_generator, but it doesn't seem possible to pass a custom task_id anywhere, nor is there uniqueness handling, I think?

  2. s3rius commented on Jan 5, 2024

    @s3rius
    Member

    Hi, @gegnew and thanks for your interest in the project. You're right. We don't handle uniqueness of ids. Ids meant to be unique only during task execution to not save calculation result twice.

    But here's what you can do. You can create a middleware that tracks ids so they are unique among tasks with specific name.

    from taskiq.abc.broker import AsyncBroker
    from taskiq.message import TaskiqMessage
    from taskiq_redis import ListQueueBroker
    from taskiq import TaskiqMiddleware
    from redis.asyncio import Redis
    
    
    class UniqueIdMiddleware(TaskiqMiddleware):
        def __init__(self, redis_url: str) -> None:
            self.pool = Redis.from_url(redis_url)
    
        def set_broker(self, broker: AsyncBroker) -> None:
            return super().set_broker(broker)
    
        async def pre_execute(self, message: TaskiqMessage) -> TaskiqMessage:
            if await self.pool.get(f"unique:{message.task_name}:{message.task_id}"):
                raise ValueError("Task is already running")
    
            await self.pool.set(f"unique:{message.task_name}:{message.task_id}", 1, ex=20)
            return message
    
    
    broker = ListQueueBroker("redis://localhost:6379").with_middlewares(
        UniqueIdMiddleware("redis://localhost:6379")
    )
    
    
    @broker.task
    async def my_task():
        print("I'm a task")
    
    
    async def main():
        await my_task.kicker().with_task_id("2").kiq()
        await my_task.kicker().with_task_id("2").kiq()
    
    
    if __name__ == "__main__":
        import asyncio
    
        asyncio.run(main())

    I won't implement this functionality on broker's behalf, because it doesn't meant to be like this. If you need it, you can easily implement it on your own.

  3. osttra-o-rotel commented on Jul 24, 2024

    @osttra-o-rotel

    Hej,

    When I tried this method, every time the ValueError gets thrown, it crashes the task runner, and it doesn't acnkowledge the message in the queue. So when the task worker restarts, I get the same task message again.

    What would be the proper way to drop the message from the queue?

  4. s3rius commented on Jul 24, 2024

    @s3rius
    Member

    Okay, I will take a look at it a bit later. Will try to come up with the fix or a new way of calling middlewares.

  5. osttra-o-rotel commented on Sep 30, 2024

    @osttra-o-rotel

    Hej s3rius!

    Any news on how to skip tasks from running from the middleware? I'm thinking changing/adding methods to the middleware tha tallows messages to be skipped, or raising a specific Skip Exception that get's handled differently where the middleware is called?

  6. tuokri commented on Jan 8, 2025

    @tuokri
    @broker.task(schedule=[{"cron": "*/1 * * * *"}])
    async def arbitrarily_long_task() -> None:
        # Depending on external factors, this task could last from a few minutes up to to ~30.
        # May start an arbitrary number of other tasks. 
        if await task_already_running():
            return
        
        # Long running work...

    I have a long running task that frequently checks whether certain external conditions are met and potentially launches a series of other long running tasks.
    The task should only run if another instance of the task is not already running.

    Currently I am setting a Redis key when the task starts and checking the existence of said key every time with task_already_running to prevent multiple instances of the task running in parallel.

    A built-in way of doing this would be useful.

  7. magnus-kr commented on Feb 5, 2026

    @magnus-kr

    it would be very useful to have some sort of idempotency when queueing tasks, and it'd even help with scheduled tasks, by having multiple instances of the scheduler running we don't have a single point of failure and don't risk the same task being executed more than once

  8. epicwhale commented on Feb 22, 2026

    @epicwhale

    it would be very useful to have some sort of idempotency when queueing tasks, and it'd even help with scheduled tasks, by having multiple instances of the scheduler running we don't have a single point of failure and don't risk the same task being executed more than once

    cannot agree more! this would be a great addition to the core library to simplify the developer experience of using it.

  9. d3vyce commented on Jun 1, 2026

    @d3vyce

    Having had the same problem, I created Taskiq Deduplication. It's a Redis-backed deduplication middleware for Taskiq that prevents duplicate tasks from being queued or executed concurrently.
    Give it a try, and don't hesitate to open an issue if it doesn't solve your problem or if a feature is missing!

  10. magnus-kr commented on Jul 28, 2026

    @magnus-kr

    is there interest in a generic way for pre_send to be able to drop tasks?
    would make things like idempotency viable, but could be used for other types of validation in pre_send as well

    i'll gladly create a PR for it, could be as simple as a SkipSendError exception and then catching it in the middleware execution loop inside kiq

  11. cgbeutler commented on Sep 29, 2026

    @cgbeutler

    @magnus-kr , I would be interested in that. That would help make deduplication much cleaner. The proposed solutions above create a lot of log noise.

  12. d3vyce commented on Sep 29, 2026

    @d3vyce

    @magnus-kr @cgbeutler I've run into the same limitation in taskiq-deduplication and had a SkipSendError implementation ready, now opened as #690. It works as proposed above: raise it from pre_send and kiq() drops the task without raising. It can also carry a task_id, so a deduplicated caller gets a handle to the task that's already queued.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions