Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
85 changes: 85 additions & 0 deletions docs/integrations/taskiq.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
# Usage with Taskiq

There is no Taskiq bootstrapper: a worker uses `FreeBootstrapper`, plus Taskiq's own OpenTelemetry
instrumentation for task spans.

## 1. Install `lite-bootstrap[free-all]` and `taskiq[opentelemetry]`:

=== "uv"

```bash
uv add "lite-bootstrap[free-all]" "taskiq[opentelemetry]"
```

=== "pip"

```bash
pip install "lite-bootstrap[free-all]" "taskiq[opentelemetry]"
```

=== "poetry"

```bash
poetry add "lite-bootstrap[free-all]" "taskiq[opentelemetry]"
```

On free-threaded CPython, install the instrument extras you need instead of `free-all`, with
`otl-http` in place of `otl`, e.g. `lite-bootstrap[sentry,logging,otl-http]`. Read more about
available extras [here](../introduction/installation.md).

## 2. Bootstrap in the worker:

```python
from taskiq import TaskiqEvents, TaskiqState
from taskiq.instrumentation import TaskiqInstrumentor

from lite_bootstrap import FreeBootstrapper, FreeConfig


broker = ... # your broker


@broker.on_event(TaskiqEvents.WORKER_STARTUP)
async def bootstrap_observability(state: TaskiqState) -> None:
state.bootstrapper = FreeBootstrapper(
FreeConfig(
service_name="my-worker",
opentelemetry_endpoint="otl",
sentry_dsn="https://testdsn@localhost/1",
)
)
state.bootstrapper.bootstrap()
TaskiqInstrumentor().instrument_broker(broker)


@broker.on_event(TaskiqEvents.WORKER_SHUTDOWN)
async def teardown_observability(state: TaskiqState) -> None:
state.bootstrapper.teardown()
```

`instrument_broker` adds Taskiq's `OpenTelemetryMiddleware` to the broker you already built, and its
spans go to the `TracerProvider` lite-bootstrap has just installed. `TaskiqInstrumentor().instrument()`
is not a substitute: it patches `AsyncBroker.__init__`, so it only reaches brokers built after it runs.

Read more about available configuration options [here](../introduction/configuration.md).

## Client and worker processes

One broker object serves both the processes that send tasks and the `taskiq worker` processes that
run them. The handlers above fire on the worker events only, so a client process is left alone: it
belongs to its own framework's bootstrapper, e.g. `FastAPIBootstrapper` in the web app that calls
`.kiq()`. Bootstrapping when the broker's module is imported would run in that process too, and set
up logging and tracing a second time.

## Metrics

lite-bootstrap has no Taskiq metrics. Taskiq's own `taskiq.middlewares.PrometheusMiddleware`
(`taskiq[metrics]`) provides them, with two side effects to know about: its constructor sets
`PROMETHEUS_MULTIPROC_DIR` for the whole process, and it serves metrics from its own HTTP server on
`server_port` (default `9000`) rather than on `prometheus_metrics_path`.

## Sentry

`sentry-python` has no Taskiq integration, and none is needed for failing tasks: Taskiq's receiver
logs a task's exception with `logger.exception`, which the `LoggingIntegration` lite-bootstrap adds
reports as an error event. Spans and breadcrumbs follow the Sentry configuration as usual.
1 change: 1 addition & 0 deletions mkdocs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ nav:
- FastStream: integrations/faststream.md
- FastAPI: integrations/fastapi.md
- FastMCP: integrations/fastmcp.md
- Taskiq: integrations/taskiq.md
- Free bootstrapper: integrations/free.md
theme:
name: material
Expand Down
2 changes: 2 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,8 @@ dev = [
# conftest imports structlog.typing at runtime, which the library does not. The suite's
# floor is higher than the one the `logging` extra publishes; this is where it belongs.
"structlog>=22.2",
# for the Taskiq recipe test only; Taskiq is not an extra.
"taskiq[opentelemetry]",
]
lint = [
"ty",
Expand Down
52 changes: 52 additions & 0 deletions tests/test_taskiq_recipe.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from taskiq import InMemoryBroker, TaskiqEvents, TaskiqState
from taskiq.instrumentation import TaskiqInstrumentor
from taskiq.middlewares.opentelemetry_middleware import OpenTelemetryMiddleware

from lite_bootstrap import FreeBootstrapper, FreeConfig


def _register_recipe(broker: InMemoryBroker, config: FreeConfig) -> None:
@broker.on_event(TaskiqEvents.WORKER_STARTUP)
async def bootstrap_observability(state: TaskiqState) -> None:
state.bootstrapper = FreeBootstrapper(config)
state.bootstrapper.bootstrap()
TaskiqInstrumentor().instrument_broker(broker)

@broker.on_event(TaskiqEvents.WORKER_SHUTDOWN)
async def teardown_observability(state: TaskiqState) -> None:
state.bootstrapper.teardown()


async def test_taskiq_recipe_traces_tasks_and_tears_down() -> None:
broker = InMemoryBroker()

@broker.task
async def add(left: int, right: int) -> int:
return left + right

_register_recipe(broker, FreeConfig(opentelemetry_log_traces=True, logging_buffer_capacity=0))

await broker.startup()
try:
bootstrapper = broker.state.bootstrapper
assert bootstrapper.is_bootstrapped
assert any(isinstance(middleware, OpenTelemetryMiddleware) for middleware in broker.middlewares)

tracer_provider = trace.get_tracer_provider()
assert isinstance(tracer_provider, TracerProvider)
exporter = InMemorySpanExporter()
tracer_provider.add_span_processor(SimpleSpanProcessor(exporter))

task = await add.kiq(1, 2)
result = await task.wait_result()
assert result.return_value == 1 + 2
span_kinds = {span.kind for span in exporter.get_finished_spans() if span.name.endswith(add.task_name)}
assert span_kinds == {trace.SpanKind.PRODUCER, trace.SpanKind.CONSUMER}
finally:
await broker.shutdown()

assert not bootstrapper.is_bootstrapped
Loading