Tasks
Sentry's task platform is designed to scale horizontally to enable high-throughput processing. The task platform is composed of a few components:
flowchart
Sentry -- produce task activation --> k[(Kafka)]
k -- consume messages --> Taskbroker
Worker -- grpc GetTask --> Taskbroker
Worker -- execute task --> Worker
Worker -- grpc SetTaskStatus --> Taskbroker
flowchart
Sentry -- produce task activation --> k[(Kafka)]
k -- consume messages --> Taskbroker
Worker -- grpc GetTask --> Taskbroker
Worker -- execute task --> Worker
Worker -- grpc SetTaskStatus --> Taskbroker
Brokers and workers are paired together to create 'processing pools' for tasks. Brokers and workers can be scaled horizontally to increase parallelism.
By default, self-hosted installs come with a single broker & worker replica. You can increase processing capacity by adding more concurrency to the single worker (via the --concurrency option on the worker), or by adding additional worker, and broker replicas. It is not recommended to go above 24 worker replicas per broker as broker performance can degrade with higher worker counts.
If your deployment requires additional processing capacity, you can add additional broker replicas and use CLI options to inform the workers of the broker addresses:
sentry run taskworker --rpc-host-list=sentry-broker-default-0:50051,sentry-broker-default-1:50051
sentry run taskworker --rpc-host-list=sentry-broker-default-0:50051,sentry-broker-default-1:50051
Workers use client-side loadbalancing to distribute load across the brokers they have been assigned to.
If you only want to scale taskworker replicas and you have a single taskbroker, you can easily increase it (up to 24 replicas) by modifying your docker-compose.override.yml file:
services:
taskworker:
deploy:
replicas: 5
services:
taskworker:
deploy:
replicas: 5
This will spawn 5 replicas of the taskworker service. Bear in mind that although the processing capacity is increased, you will need to monitor and scale your system resources (CPU and RAM) accordingly.
If you have reached the maximum number of taskworker replicas, you can scale your taskbroker replicas. First, you need to increase the number of partitions on "taskworker" topic in Kafka, the number of partitions should be evenly divisible by the number of taskbroker replicas that you're planning to scale to. Then due to the limited capability of Docker Compose, you will need to manually scale your taskbroker replicas. You can do this by adding more container declarations to your docker-compose.override.yml file. Each replica reuses the same taskbroker/config.yml that ships with self-hosted (the one the default taskbroker service already mounts), so they consume the same topics and join the same consumer groups — only the SQLite volume differs per replica:
services:
taskbroker-beta:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.yml
volumes:
- sentry-taskbroker-beta:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
taskbroker-charlie:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.yml
volumes:
- sentry-taskbroker-charlie:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
volumes:
sentry-taskbroker-beta: {}
sentry-taskbroker-charlie: {}
services:
taskbroker-beta:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.yml
volumes:
- sentry-taskbroker-beta:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
taskbroker-charlie:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.yml
volumes:
- sentry-taskbroker-charlie:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
volumes:
sentry-taskbroker-beta: {}
sentry-taskbroker-charlie: {}
Note that each taskbroker replica needs their own SQLite database per replica, to prevent issues contention and locks between replicas.
Finally, you need to modify taskworker command to have rpc-host-list pointing to the new brokers:
taskworker:
<<: *sentry_defaults
command: run taskworker --concurrency=4 --rpc-host-list=taskbroker:50051,taskbroker-beta:50051,taskbroker-charlie:50051 --health-check-file-path=/tmp/health.txt
healthcheck:
<<: *file_healthcheck_defaults
taskworker:
<<: *sentry_defaults
command: run taskworker --concurrency=4 --rpc-host-list=taskbroker:50051,taskbroker-beta:50051,taskbroker-charlie:50051 --health-check-file-path=/tmp/health.txt
healthcheck:
<<: *file_healthcheck_defaults
If you are running on Kubernetes, you can utilize --rpc-host and --num-brokers instead if you are using StatefulSet. Otherwise, you can use --rpc-host-list if you have a different host name pattern for the brokers.
In higher throughput installations, you may also want to isolate task workloads from each other to ensure timely processing of lower volume tasks. For example, you could isolate ingestion related tasks from other work:
flowchart
Sentry -- produce tasks --> k[(Kafka)]
k -- topic-taskworker-ingest --> tb-i[Taskbroker ingest]
k -- topic-taskworker --> tb-d[Taskbroker default]
tb-i --> w-i[ingest Worker]
tb-d --> w-d[default Worker]
flowchart
Sentry -- produce tasks --> k[(Kafka)]
k -- topic-taskworker-ingest --> tb-i[Taskbroker ingest]
k -- topic-taskworker --> tb-d[Taskbroker default]
tb-i --> w-i[ingest Worker]
tb-d --> w-d[default Worker]
To achieve this work separation we need to make a few changes:
- Provision any additional topics. Topic names need to come from one of the predefined topics in
src/sentry/conf/types/kafka_definition.py. By default, any topics will automatically be created during./install.shprocess. - Deploy the additional broker replicas. Give each broker its own config file declaring the topic it consumes under
kafka_topics, and point the broker at it with the-cflag. - Deploy additional workers that use the new brokers in their
rpc-host-listCLI flag. - Find the list of namespaces you want to shift to the new topic. The list of task namespaces can be found in the
sentry.taskworker.namespacesmodule. - Update task routing option, defining the namespace -> topic mappings. e.g.Copied
# in sentry/config.yml taskworker.route.overrides: "ingest.errors": "taskworker-ingest" "ingest.transactions": "taskworker-ingest"# in sentry/config.yml taskworker.route.overrides: "ingest.errors": "taskworker-ingest" "ingest.transactions": "taskworker-ingest"
Having separate ingest taskbroker and taskworker is useful for high-throughput installations, therefore you can receive timely alerts and not have to wait for ingest-related tasks to finish. As an implementation of the above steps, you need to create a config file for the ingest broker that declares the taskworker-ingest topic. Create taskbroker/config.ingest.yml:
kafka_deadletter_topic: taskworker-ingest-dlq
kafka_topics:
taskworker-ingest:
cluster: default
consumer_group: taskworker-ingest
taskworker-ingest-dlq:
cluster: default
consumer_group: taskworker-ingest
produce_only: true
kafka_deadletter_topic: taskworker-ingest-dlq
kafka_topics:
taskworker-ingest:
cluster: default
consumer_group: taskworker-ingest
taskworker-ingest-dlq:
cluster: default
consumer_group: taskworker-ingest
produce_only: true
Then add a few new containers on your docker-compose.override.yml file:
# Copy `x-sentry_defaults` and `file_healthcheck_defaults` section from
# `docker-compose.yml` to `docker-compose.override.yml` first. Put it on the
# top of the file.
services:
taskbroker-ingest:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations-ingest.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.ingest.yml
volumes:
- sentry-taskbroker-ingest:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
taskworker-ingest:
<<: *sentry_defaults
command: run taskworker --concurrency=4 --rpc-host-list=taskbroker-ingest:50051 --health-check-file-path=/tmp/health.txt
healthcheck:
<<: *file_healthcheck_defaults
volumes:
sentry-taskbroker-ingest: {}
# Copy `x-sentry_defaults` and `file_healthcheck_defaults` section from
# `docker-compose.yml` to `docker-compose.override.yml` first. Put it on the
# top of the file.
services:
taskbroker-ingest:
restart: "unless-stopped"
image: "$TASKBROKER_IMAGE"
environment:
TASKBROKER_KAFKA_CLUSTERS__DEFAULT__ADDRESS: "kafka:9092"
TASKBROKER_DB_PATH: "/opt/sqlite/taskbroker-activations-ingest.sqlite"
TASKBROKER_STATSD_ADDR: ${STATSD_ADDR:-127.0.0.1:8125}
command: /opt/taskbroker -c /etc/taskbroker/config.ingest.yml
volumes:
- sentry-taskbroker-ingest:/opt/sqlite
- type: bind
read_only: true
source: ./taskbroker
target: /etc/taskbroker
depends_on:
- kafka
taskworker-ingest:
<<: *sentry_defaults
command: run taskworker --concurrency=4 --rpc-host-list=taskbroker-ingest:50051 --health-check-file-path=/tmp/health.txt
healthcheck:
<<: *file_healthcheck_defaults
volumes:
sentry-taskbroker-ingest: {}
On your sentry/config.yml file, you need to append the following to the bottom of the file:
taskworker.route.overrides:
"ingest.errors": "taskworker-ingest"
"ingest.transactions": "taskworker-ingest"
"ingest.profiling": "taskworker-ingest"
"ingest.attachments": "taskworker-ingest"
"ingest.errors.postprocess": "taskworker-ingest"
taskworker.route.overrides:
"ingest.errors": "taskworker-ingest"
"ingest.transactions": "taskworker-ingest"
"ingest.profiling": "taskworker-ingest"
"ingest.attachments": "taskworker-ingest"
"ingest.errors.postprocess": "taskworker-ingest"
Any other tasks that are not defined on the routes override above will be handled by the default taskbroker and taskworker service.
Our documentation is open source and available on GitHub. Your contributions are welcome, whether fixing a typo (drat!) or suggesting an update ("yeah, this would be better").