---
title: "Transactional ack and push"
description: "The acknowledgement of one stage and the push of the next, in one PostgreSQL transaction."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Transactional ack and push

The previous tutorial handed work between queues in two steps. In that order a crash between
them duplicates work; in the other order it loses it. There is no arrangement of two calls that
avoids both.

A transaction removes the choice: the ack of the input and the push of the output commit together
or not at all. The commit also fails if the lease has expired, which is what stops a slow consumer
from acknowledging work the broker has already handed to somebody else.

Two details make it work in the loop. `auto_ack(False)` stops `consume()` from acknowledging
behind your back, since the ack now belongs to the transaction. And the ack has to name its
consumer group, spelled `consumer_group` here rather than in the wire's camelCase: the transaction
does not read it off the message, and an ack without it commits the wrong cursor.

This is atomic handoff, not end-to-end exactly-once delivery: a lost response plus a retry still
duplicates the pushes unless their transaction ids are deterministic. What it buys is that your
pipeline's state can never disagree with itself.

```python title="examples/tutorials/py/03_transaction_ack_push.py"
#
# Tutorial 3 of 5: acknowledge and push in one transaction.
#
# Tutorial 2 handed work from one queue to the next in two steps: push the
# derived message, then let the loop acknowledge the source. Between those two
# steps a crash duplicates work, and in the other order it loses work.
#
# A Queen transaction closes that window: the acknowledgement of the input and
# the push of the output are one PostgreSQL transaction. Both land or neither
# does.
#
# Run it:
#   QUEEN_URL=http://localhost:6632 python3 03_transaction_ack_push.py

import asyncio
import os
import sys
import time

from queen import Queen

QUEEN_URL = os.environ.get("QUEEN_URL", "http://localhost:6632")
RUN = f"{int(time.time() * 1000):x}"
ORDERS = f"tut-py-tx-orders-{RUN}"
INVOICES = f"tut-py-tx-invoices-{RUN}"
GROUP = "tut-py-invoicing"

INPUT = [
    {"orderId": "A-1", "customer": "acme", "total": 120.5},
    {"orderId": "B-1", "customer": "globex", "total": 88.75},
    {"orderId": "C-1", "customer": "initech", "total": 310.0},
]

CHECKS = 0


def check(condition: bool, description: str) -> None:
    """Record one verified fact, or abort the run."""
    global CHECKS
    if not condition:
        raise AssertionError(description)
    CHECKS += 1
    print(f"  ok: {description}")


async def main() -> int:
    queen = Queen(url=QUEEN_URL)
    verdict, failed = "", False

    try:
        print(f"broker {QUEEN_URL}")

        for order in INPUT:
            await queen.queue(ORDERS).partition(order["customer"]).push({"data": order})
        print(f"pushed {len(INPUT)} orders")

        print("\ninvoicing")
        invoiced = []

        async def invoice(msg):
            # One commit carries both operations. The ack names the consumer
            # group explicitly: the transaction builder does not read it off the
            # message, and an ack sent without it commits the wrong cursor. The
            # Python builder spells that context key `consumer_group`, unlike
            # the camelCase the wire uses.
            result = await (
                queen.transaction()
                .queue(INVOICES)
                .partition(msg["data"]["customer"])
                .push(
                    [
                        {
                            "data": {
                                "invoiceId": f"INV-{msg['data']['orderId']}",
                                "orderId": msg["data"]["orderId"],
                                "amount": msg["data"]["total"],
                            }
                        }
                    ]
                )
                .ack(msg, "completed", {"consumer_group": GROUP})
                .commit()
            )

            # Check the transaction, not just the absence of an exception.
            # commit() does raise on a rejected transaction, so this is belt and
            # braces against a success flag that says otherwise.
            if not result.get("success"):
                raise RuntimeError(f"transaction rejected: {result.get('error')}")

            invoiced.append(msg["data"]["orderId"])
            print(f"  {msg['data']['orderId']} -> INV-{msg['data']['orderId']}")

        # auto_ack(False) is what makes this tutorial possible: the loop must not
        # acknowledge behind your back, because the acknowledgement is part of
        # the transaction above.
        await (
            queen.queue(ORDERS)
            .group(GROUP)
            .subscription_mode("all")
            .each()
            .auto_ack(False)
            .limit(len(INPUT))
            .idle_millis(5000)
            .consume(invoice)
        )

        check(len(invoiced) == len(INPUT), "every order was invoiced once")

        # The commit fails if the lease has expired, which is what stops a slow
        # consumer from acking work the broker has already handed to someone
        # else. Nothing to assert here: the check above is that assertion, since
        # a failed commit would have raised.

        print("\nchecking the output queue")

        # The invoices went to one partition per customer, and a pop claims a
        # single partition unless you say otherwise: partitions(10) lets this one
        # call claim up to ten of them, with batch as the total budget across all
        # of them.
        invoices = await queen.queue(INVOICES).batch(10).partitions(10).wait(True).pop()

        check(len(invoices) == len(INPUT), f"{len(INPUT)} invoices exist")
        ids = sorted(m["data"]["orderId"] for m in invoices)
        check(
            ids == sorted(o["orderId"] for o in INPUT),
            "each invoice matches an order, none duplicated",
        )

        # And the input queue is committed for this group: the acks were part of
        # the same transactions that produced those invoices, so the two states
        # cannot disagree.
        leftovers = await queen.queue(ORDERS).group(GROUP).batch(10).wait(False).pop()
        check(len(leftovers) == 0, "the source queue is committed for this group")

        await queen.queue(ORDERS).delete()
        await queen.queue(INVOICES).delete()

        verdict = f"\nPASS: {CHECKS} checks"
    except Exception as err:
        verdict, failed = f"\nFAIL: {err}", True
    finally:
        await queen.close()

    # A failure goes to stderr, like the rest of the set. Flush stdout first so
    # the verdict still lands last when the two are piped into one file.
    sys.stdout.flush()
    print(verdict, file=sys.stderr if failed else sys.stdout)
    return 1 if failed else 0


if __name__ == "__main__":
    sys.exit(asyncio.run(main()))
```

## Run it

Against a broker from [the quickstart](/start/quickstart), from the root of the repository, with
the in-tree client on `PYTHONPATH`:

```bash
PYTHONPATH=clients/client-py python3 examples/tutorials/py/03_transaction_ack_push.py
```

The program checks its own outcome and exits non-zero if a check fails. Every tutorial in this
section runs in the repository's own suite: `examples/tutorials/run.sh py`.

Next: [Replay](/use/python-client/replay).

Source: https://queenmq.com/use/python-client/transaction-ack-push/index.mdx
