Skip to main content

Rejected is a AMQP consumer daemon and message processing framework. It allows for rapid development of message processing consumers by handling all of the core functionality of communicating with RabbitMQ and management of consumer processes.

Rejected runs as a master process with multiple consumer configurations that are each run it an isolated process. It has the ability to collect statistical data from the consumer processes and report on it.

Rejected supports Python 2.7 and 3.4+.

Version Status Coverage License

Features

  • Automatic exception handling including connection management and consumer restarting

  • Smart consumer classes that can automatically decode and deserialize message bodies based upon message headers

  • Metrics logging and submission to statsd and InfluxDB

  • Built-in profiling of consumer code

  • Ability to write asynchronous code in consumers allowing for parallel communication with external resources

Documentation

https://rejected.readthedocs.io

Example Consumers

from rejected import consumer
import logging

LOGGER = logging.getLogger(__name__)


class Test(consumer.Consumer):

    def process(self, message):
        LOGGER.debug('In Test.process: %s' % message.body)

Async Consumer

To make a consumer async, you can decorate the Consumer.prepare and Consumer.process methods using Tornado’s @gen.coroutine. Asynchronous consumers do not allow for concurrent processing multiple messages in the same process, but rather allow you to use asynchronous clients like Tornado’s AsyncHTTPClient and the Queries PostgreSQL library to perform parallel tasks using coroutines when processing a single message.

import logging

from rejected import consumer

from tornado import gen
from tornado import httpclient


class AsyncExampleConsumer(consumer.Consumer):

    @gen.coroutine
    def process(self):
        LOGGER.debug('Message: %r', self.body)
        http_client = httpclient.AsyncHTTPClient()
        results = yield [http_client.fetch('http://www.github.com'),
                         http_client.fetch('http://www.reddit.com')]
        LOGGER.info('Length: %r', [len(r.body) for r in results])

Example Configuration

%YAML 1.2
---
Application:
  poll_interval: 10.0
  stats:
    log: True
    influxdb:
      enabled: True
      scheme: http
      host: localhost
      port: 8086
      user: username
      password: password
      database: dbname
    statsd:
      enabled: True
      host: localhost
      port: 8125
      prefix: applications.rejected
  Connections:
    rabbitmq:
      host: localhost
      port: 5672
      user: guest
      pass: guest
      ssl: False
      vhost: /
      heartbeat_interval: 300
  Consumers:
    example:
      consumer: rejected.example.Consumer
      sentry_dsn: https://[YOUR-SENTRY-DSN]
      connections:
        - name: rabbitmq1
          consume: True
      drop_exchange: dlxname
      qty: 2
      queue: generated_messages
      qos_prefetch: 100
      ack: True
      max_errors: 100
      config:
        foo: True
        bar: baz

Daemon:
  user: rejected
  group: daemon
  pidfile: /var/run/rejected/example.%(pid)s.pid

Logging:
  version: 1
  formatters:
    verbose:
      format: "%(levelname) -10s %(asctime)s %(process)-6d %(processName) -25s %(name) -20s %(funcName) -25s: %(message)s"
      datefmt: "%Y-%m-%d %H:%M:%S"
    verbose_correlation:
      format: "%(levelname) -10s %(asctime)s %(process)-6d %(processName) -25s %(name) -20s %(funcName) -25s: %(message)s {CID %(correlation_id)s}"
      datefmt: "%Y-%m-%d %H:%M:%S"
    syslog:
      format: "%(levelname)s <PID %(process)d:%(processName)s> %(name)s.%(funcName)s: %(message)s"
    syslog_correlation:
      format: "%(levelname)s <PID %(process)d:%(processName)s> %(name)s.%(funcName)s: %(message)s {CID %(correlation_id)s)"
  filters:
    correlation:
      '()': rejected.log.CorrelationFilter
      'exists': True
    no_correlation:
      '()': rejected.log.CorrelationFilter
      'exists': False
  handlers:
    console:
      class: logging.StreamHandler
      formatter: verbose
      debug_only: false
      filters: [no_correlation]
    console_correlation:
      class: logging.StreamHandler
      formatter: verbose_correlation
      debug_only: false
      filters: [correlation]
    syslog:
      class: logging.handlers.SysLogHandler
      facility: daemon
      address: /var/run/syslog
      formatter: syslog
      filters: [no_correlation]
    syslog_correlation:
      class: logging.handlers.SysLogHandler
      facility: daemon
      address: /var/run/syslog
      formatter: syslog
      filters: [correlation]
  loggers:
    helper:
      level: INFO
      propagate: true
      handlers: [console, console_correlation, syslog, syslog_correlation]
    rejected:
      level: INFO
      propagate: true
      handlers: [console, console_correlation, syslog, syslog_correlation]
    tornado:
      level: INFO
      propagate: true
      handlers: [console, console_correlation, syslog, syslog_correlation]
  disable_existing_loggers: true
  incremental: false

Version History

Available at https://rejected.readthedocs.org/en/latest/history.html

Release files for rejected 3.23.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for rejected 3.23.1
File Size Uploaded
rejected-3.23.1.tar.gz 53.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for rejected 3.23.1
File Interpreter ABI Platform
rejected-3.23.1-py2.py3-none-any.whl Python 3, Python 2 none any Details

Total release size: 101.8 kB

Release files / rejected-3.23.1.tar.gz

Download URL rejected-3.23.1.tar.gz
Size 53.7 kB
Tags Source
SHA-256 checksum
How to use checksums
d56047a549037463b9a3a64e8c75a3f0e08c3209ece3a0ed53b8e4131b4d3678
BLAKE2b-256 checksum
How to use checksums
6918353da913990d9f60352f8a888a2f1b08c6efd6a0c925ef251faabcb013fc
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.9.16

Release files / rejected-3.23.1-py2.py3-none-any.whl

Download URL rejected-3.23.1-py2.py3-none-any.whl
Size 48.2 kB
Tags Python 2 Python 3
SHA-256 checksum
How to use checksums
33aaa373cb492c689444d526360a1e4c9e78d7971e1e65bc60562a5dfc2523fa
BLAKE2b-256 checksum
How to use checksums
56a17d9ed7c2c1c945726eb0d5c4a1a385b5f8a6a0e215fb0ec62163212190c3
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.9.16

Release history Release notifications | RSS feed

This release

3.23.1 This release

2 release files

3.23.0

2 release files

3.22.3

2 release files

3.22.2

2 release files

3.22.1

2 release files

3.22.0

1 release file

3.21.1

1 release file

3.20.9

2 release files

3.20.8

2 release files

3.20.7

2 release files

3.20.6

2 release files

3.20.5

2 release files

3.20.1

1 release file

3.20.0

2 release files

3.19.9

2 release files

3.19.8

2 release files

3.19.7

2 release files

3.19.4

2 release files

3.19.3

2 release files

3.19.2

2 release files

3.19.1

2 release files

3.19.0

2 release files

3.18.9

2 release files

3.18.8

2 release files

3.18.5

2 release files

3.18.4

2 release files

3.18.1

2 release files

3.17.4

2 release files

3.17.3

2 release files

3.17.2

2 release files

3.17.1

2 release files

3.17.0

2 release files

3.16.7

2 release files

3.16.6

2 release files

3.16.5

2 release files

3.16.4

2 release files

3.16.1

2 release files

3.16.0

2 release files

3.15.1

2 release files

3.15.0

2 release files

3.14.0

2 release files

3.12.1

2 release files

3.11.1

2 release files

3.11.0

2 release files

3.10.0

2 release files

3.9.4

2 release files

3.9.3

2 release files

3.9.2

2 release files

3.9.1

2 release files

3.9.0

2 release files

3.8.0

2 release files

3.7.2

2 release files

3.7.1

2 release files

3.7.0

2 release files

3.6.4

2 release files

3.6.3

2 release files

3.6.2

2 release files

3.6.1

2 release files

3.6.0

2 release files

3.5.1

2 release files

3.5.0

2 release files

3.4.6

2 release files

3.4.5

2 release files

3.4.4

2 release files

3.4.3

2 release files

3.4.2

2 release files

3.4.1

2 release files

3.4.0

2 release files

3.3.0

1 release file

3.2.19

1 release file

3.2.18

1 release file

3.2.17

1 release file

3.2.16

1 release file

3.2.15

1 release file

3.2.14

1 release file

3.2.13

3.2.11

1 release file

3.2.10

1 release file

3.2.9

1 release file

3.2.8

1 release file

3.2.6

1 release file

3.2.5

1 release file

3.2.4

1 release file

3.2.3

1 release file

3.2.2

1 release file

3.2.1

1 release file

3.2.0

1 release file

3.1.2

1 release file

3.0.30

1 release file

3.0.29

1 release file

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page