AMQP 1.0 Source Connectors

Type

Source

Core Class

io.axual.connect.plugins.amqp.source.AmqpSourceConnector

OAuth2 Class

io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector

Target System

Messaging & Streaming (AMQP 1.0)

Maintainer

Axual

License

Proprietary (client-only)

Project

Proprietary. Source code is not publicly accessible.

Download

Contact Axual Support to obtain the connector library.

This page documents version 1.1.0. Newer versions should be compatible unless there are breaking changes, but field names or default values may differ. If you notice discrepancies, please contact Axual Support.

Description

The AMQP 1.0 Source Connectors read messages from AMQP 1.0 queues or topics and publish them as records to Kafka. They are tested with ActiveMQ Artemis and RabbitMQ (with AMQP 1.0 plugin).

Two connector variants are available:

  • Core AMQP Source Connector (io.axual.connect.plugins.amqp.source.AmqpSourceConnector) - supports anonymous or username/password (PLAIN) authentication.

  • OAuth2 AMQP Source Connector (io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector) - supports token-based OAuth2 authentication.

Both connectors share the same AMQP consumption behaviour (queue/topic address, QoS, link durability, etc.). The OAuth2 variant replaces basic authentication with token acquisition and SASL/OAuth2.

Connector variants

Select a connector based on the authentication offered by your AMQP broker or gateway:

Authentication scenario Connector to use

No authentication (anonymous)

Core connector (io.axual.connect.plugins.amqp.source.AmqpSourceConnector).
Leave amqp.auth.username and amqp.auth.password unset to log in anonymously, or set amqp.auth.enabled=false to skip the SASL exchange altogether.

Username and password (SASL PLAIN)

Core connector. Set amqp.auth.username and amqp.auth.password.

OAuth2 (client credentials or private key JWT)

OAuth2 connector (io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector). Configure amqp.oauth. properties.
The connector sets amqp.auth.
itself from the acquired token, so configuring them here has no effect.

OAuth2 method selection (amqp.oauth.method):

Method Description

client_secret_basic

Client ID and secret sent via HTTP Basic.

client_secret_post

Client ID and secret sent in the token request body.

private_key_jwt

Signed client assertion using a private key.
Provide the key as JWK JSON (amqp.oauth.private.key.jwk.json) or PEM (amqp.oauth.private.key.pem.content).

Features

  • Read messages from AMQP 1.0 queues or topics and publish them as records to Kafka

  • Two variants: Core (anonymous/PLAIN) and OAuth2 (client credentials or private key JWT)

  • Configurable Quality of Service: AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE

  • Durable consumer support - AMQP link persists across disconnections so no messages are lost on restart

  • Flexible key extraction chain: MESSAGE_ID, CORRELATION_ID, SUBJECT, UUID, or NULL

  • Message bodies from an AMQP Data section or Value section, published Base64 encoded or decoded as JSON (amqp.data.mode)

  • AMQP message properties forwarded as Kafka headers with msg. prefix

  • AMQP application properties forwarded as Kafka headers with app. prefix

  • Header values stored as bytes or as readable text (amqp.header.value.format)

  • Configurable handling of a message that cannot be converted: drop it, hand it back to the broker for dead-lettering, or fail the task (amqp.error.handling)

  • Ships with a transformation that decodes a Base64 encoded field inside a JSON body, see DecodeBase64JsonField

  • Tested with ActiveMQ Artemis and RabbitMQ (AMQP 1.0 plugin)

When to Use

  • You need to bridge an AMQP 1.0 broker (RabbitMQ, ActiveMQ Artemis) to Kafka.

  • Your broker requires OAuth2 token-based authentication.

  • You need durable subscriptions to ensure no messages are lost across connector restarts.

  • You want configurable QoS levels for message delivery guarantees.

When NOT to Use

  • Your broker does not support AMQP 1.0 - AMQP 0.9.x (classic RabbitMQ default) requires a different connector.

  • You need schema-aware (structured) Kafka records - this connector always emits plain string keys and values.

  • You need to enrich messages, join them with other data, or reshape them before publishing to Kafka - use KSML or Kafka Streams after ingestion instead. The connector applies only the Single Message Transformations you configure, such as the bundled DecodeBase64JsonField.

Installation

The AMQP 1.0 Source Connectors are maintained by Axual in a private repository and are not published to a public artifact repository.

Contact the Axual Support team to obtain the connector library.

For how a plugin reaches the workers, see Configure plugins.

Configuration

Core AMQP Source Connector

The following options are available for the Core connector (io.axual.connect.plugins.amqp.source.AmqpSourceConnector).

Property Description Type Default Valid values Importance

amqp.host

The hostname where the AMQP server can be reached.

string

non-empty string

high

amqp.port

The port number where the AMQP server listens.

int

5672

[1,…​]

medium

amqp.tls.enabled

Enable TLS connection to the AMQP server.

boolean

false

medium

amqp.auth.enabled

When set to false, no authentication is attempted. Default is true.

boolean

true

high

amqp.auth.username

Username for connecting to the AMQP server. Uses anonymous authentication if not set.

string

null

high

amqp.auth.password

Password for connecting to the AMQP server.

password

null

high

amqp.source

The AMQP queue or topic address to consume messages from. This corresponds to the queue name in the broker.

string

non-empty string

high

amqp.consumer.qos

Quality of Service level: AT_MOST_ONCE (fire and forget), AT_LEAST_ONCE (guaranteed, may duplicate), EXACTLY_ONCE (requires broker support).

string

AT_LEAST_ONCE

AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE

medium

amqp.consumer.durable

Use a durable consumer. The AMQP link remains in place after disconnection so messages are not lost. Essential for persistent topic subscriptions.

boolean

false

medium

amqp.consumer.link.name

Explicit link name. Overrides auto-generation. When multiple tasks share the same link name, messages are load-balanced between them.

string

null

low

amqp.consumer.link.name.prefix

Prefix for auto-generated link names. The actual link name will be {prefix}-{task-id}. When unset, the connector name is used as the prefix, so the name becomes {connector-name}-{task-id}.

string

null

medium

amqp.consumer.link.credit

Number of messages prefetched and buffered locally before requiring acknowledgments.

int

100

[1,…​]

medium

amqp.consumer.receive.timeout

Maximum time in milliseconds to wait for a message when polling.

long

1000

[100,…​]

low

amqp.consumer.close.timeout

Maximum time in milliseconds to wait for graceful consumer shutdown.

long

5000

[1000,…​]

low

amqp.consumer.incoming.window

Session-level incoming window size for flow control.

int

2048

[1,…​]

low

amqp.consumer.outgoing.window

Session-level outgoing window size for flow control.

int

2048

[1,…​]

low

amqp.consumer.auto.acknowledge

Automatically acknowledge messages after successful processing.

boolean

false

medium

amqp.consumer.reject.on.error

On processing error: if true, reject the message (no redelivery); if false, release it (can be redelivered).

boolean

false

medium

amqp.idle.timeout

Maximum idle time in milliseconds before the connection is considered dead.

long

2147483647

[1,…​]

medium

amqp.max.frame.size

Maximum size in bytes for individual AMQP frames.

long

2147483647

[1,…​]

medium

amqp.input.buffer.size

Network buffer size in bytes for receiving frames.

int

131072

[1,…​]

medium

amqp.input.buffer.extend.size

Extension size in bytes when the input buffer is insufficient to receive a frame.

int

65536

[1,…​]

medium

amqp.output.buffer.size

Network buffer size in bytes for sending frames.

int

131072

[1,…​]

medium

amqp.output.buffer.extend.size

Extension size in bytes when the output buffer is insufficient to send a frame.

int

65536

[1,…​]

medium

amqp.container.id.prefix

Prefix for auto-generated AMQP container IDs. The actual ID will be {prefix}-{task-id}. When unset, the connector name is used as the prefix, so the ID becomes {connector-name}-{task-id}.

string

null

medium

amqp.container.id

Explicit container ID. Overrides auto-generation.
WARNING: Using the same ID across multiple tasks may cause broker tracking issues.

string

null

low

amqp.key.extractor

Ordered list of key extraction strategies. Each is tried in turn until a key is found.
Options (case insensitive): MESSAGE_ID, CORRELATION_ID, SUBJECT, UUID, NULL.

list

MESSAGE_ID,SUBJECT,NULL

high

amqp.data.mode

How an encoded message body becomes the Kafka record value. See Message body handling.
BASE64 Base64 encodes the Data section bytes and leaves Value bodies rendered as in release 1.0.0.
JSON decodes the body bytes as UTF-8 and publishes the JSON text unchanged. It applies to a Data section and to a Value section holding binary; text in a Value section is published unchanged, because text needs no decoding.
BASE64_JSON reads the body as Base64 text, decodes that first, and then behaves as JSON. It also applies to a Value section holding text.
In both decoding modes, a body that does not hold exactly one valid JSON document is a conversion failure, handled according to amqp.error.handling.

string

BASE64

BASE64, JSON, BASE64_JSON

high

amqp.header.value.format

How the msg. and app. record headers are stored. See Header values.
BYTES stores them as UTF-8 bytes, which the default header converter renders Base64 encoded.
STRING stores them as strings, which the default header converter renders as readable text.

string

BYTES

BYTES, STRING

medium

amqp.error.handling

What happens to a message that could not be converted. See Conversion error handling.
When left empty, the deprecated amqp.unsupported.type.ignore decides: true means DROP, false means FAIL.

string

DROP, BROKER_REJECT, FAIL

medium

amqp.unsupported.type.ignore

Deprecated, superseded by amqp.error.handling, which is used instead whenever it is set. When true, a message that cannot be converted is not written to Kafka; when false, the task fails.

boolean

true

medium

amqp.unsupported.type.log

Log a message that could not be converted.

boolean

true

medium

amqp.unsupported.type.log.name

Logger name to use when logging a message that could not be converted. Empty means the connector’s own logger.

string

low

amqp.unsupported.type.log.include.body

Include the full AMQP message in that log entry.

boolean

false

low

amqp.unsupported.type.log.include.message.id

Include the message id in that log entry. Ignored when the full message is already included.

boolean

false

low

topic

The Kafka topic to write messages to.

string

non-empty string

high

kafka.commit.timeout.ms

How long to wait for a converted message to be confirmed written to Kafka. While messages are waiting, the AMQP consumer is paused.

long

30000

[1,…​]

medium

amqp.validate.connection

Attempt a live broker connection during connector validation.

boolean

false

low

OAuth2 AMQP Source Connector

The OAuth2 connector (io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector) takes the same AMQP connection, consumer, tuning and conversion properties as the Core connector, with two exceptions. It does not define the amqp.auth.* properties, because it sets them itself from the acquired token before handing the configuration to its tasks, so any value you configure under that prefix is overwritten. It also does not define amqp.validate.connection, so the live broker check is not available on this variant. It adds the following OAuth2-specific options:

Property Description Type Default Importance

amqp.oauth.username

The AMQP username for SASL authentication. The bearer token is sent as the password, so this is the only credential you set yourself.

string

high

amqp.oauth.token.endpoint

The URL for the access token endpoint.

string

high

amqp.oauth.method

Client authentication method: client_secret_basic, client_secret_post, or private_key_jwt.

string

client_secret_basic

medium

amqp.oauth.client.id

The client identifier registered with the authorization server.

string

high

amqp.oauth.client.secret

The client secret (for client_secret_basic and client_secret_post methods).

password

null

medium

amqp.oauth.scope

Space-separated list of requested permissions (e.g., read_users write_data).

string

null

medium

amqp.oauth.private.key.jwk.json

Private key in JWK JSON format (for private_key_jwt method).

password

null

medium

amqp.oauth.private.key.pem.content

Private key in PEM encoding (alternative to JWK).

password

null

medium

amqp.oauth.private.key.pem.password

Password for an encrypted PEM private key.

password

null

low

amqp.oauth.private.key.jws.algorithm

JWS algorithm used to sign the client assertion. Default is RS256.

string

RS256

medium

amqp.oauth.private.key.claims.issuer

JWT Issuer Claim (iss). Defaults to amqp.oauth.client.id if not set.

string

null

medium

amqp.oauth.private.key.claims.audience

JWT Audience Claim (aud). Defaults to amqp.oauth.token.endpoint if not set.

string

null

medium

amqp.oauth.private.key.claims.expirationSeconds

JWT Expiration Time Claim (exp) in seconds.

long

300

medium

amqp.oauth.private.key.claims.notBeforeSeconds

JWT Not Before Claim (nbf). Seconds added to the issued-at time. Not set by default.

long

null

medium

amqp.oauth.private.key.claims.issuedAt

When true, adds an Issued At (iat) claim to the JWT.

boolean

false

medium

amqp.oauth.private.key.claims.jwtId

When true, adds a random JWT Identifier (jti) claim.

boolean

false

medium

Data conversion

The connector converts AMQP 1.0 messages to Kafka Connect records as follows:

  • Message body (payload) - always emitted as a String value using Schema.OPTIONAL_STRING_SCHEMA. AMQPNull produces a null record value. How an encoded body becomes that string is set by amqp.data.mode, see Message body handling.

  • Record key - extracted by the amqp.key.extractor chain using the same string formatting.

  • Message Properties - added as Kafka headers with the msg. prefix (e.g., msg.Subject, msg.ContentType, msg.CreationTime).

  • Application Properties - added as Kafka headers with the app. prefix. Container types (Array/List/Map) are skipped.

  • Body section - reported as the msg.amqp.type header, with the value Data or Value.

Header values are always derived as text. amqp.header.value.format sets whether they are stored as bytes or as strings, see Header values. Keys and header values keep the type specific rendering below in every mode; amqp.data.mode applies to the body only.

Type-specific string rendering for body and header values:

AMQP type String representation

AMQPBinary

"binary : " + Base64

AMQPDecimal128

"decimal128 : " + Base64

AMQPDecimal32 / AMQPDecimal64

Decimal notation, as BigDecimal.toString() renders it

AMQPByte / AMQPUnsignedByte

"0x" + 2-digit hex

AMQPTimestamp

Epoch milliseconds (long as string)

AMQPBoolean

"true" or "false" (lowercase)

AMQPChar

Single-character Unicode string

AMQPUuid

Canonical UUID text, as UUID.toString() renders it

AMQPString

The string itself

Numeric types (AMQPInt, AMQPLong, AMQPShort, AMQPFloat, AMQPDouble and their unsigned variants)

String.valueOf(value)

AMQPArray / AMQPList / AMQPMap

Not supported - results in a conversion error

Any other type

The string form the AMQP type reports for itself

Message body handling

AMQP 1.0 carries the message body in one of three sections: a Value section (spec 3.2.8), a Data section (spec 3.2.6, opaque bytes) or a Sequence section (spec 3.2.7, not supported). A producer that sends JSON may use either of the first two, and may encode it before sending.

amqp.data.mode describes the body, not the section that carries it, so the same setting works whichever section your producer uses:

Body BASE64 (default) JSON BASE64_JSON

Data section

Base64 text of the bytes

Bytes decoded as UTF-8 and published unchanged

Read as Base64 first, then decoded as UTF-8 and published unchanged

Value section holding binary

binary : <base64>

Bytes decoded as UTF-8 and published unchanged

Read as Base64 first, then decoded as UTF-8 and published unchanged

Value section holding text

Published unchanged

Published unchanged, text needs no decoding

Read as Base64, decoded and published unchanged

Value section, any other type

Rendered per its AMQP type in every mode, see the table above

In JSON and BASE64_JSON the decoded text must be exactly one well-formed JSON document. It is published exactly as the producer wrote it, never re-serialised, so field order, spacing and the exact spelling of numbers are preserved. A body that is not valid UTF-8, not valid Base64, not valid JSON, or that concatenates two documents, is a conversion failure and is handled according to amqp.error.handling.

The record value stays a string in every mode, so value.converter remains org.apache.kafka.connect.storage.StringConverter. Consumers of a topic filled in JSON mode read the JSON directly instead of Base64 decoding it first.

A body split over several Data sections is joined: in BASE64 each section is encoded separately and the results are joined with a comma, as in release 1.0.0; in the decoding modes the bytes are joined first and then decoded once, as the AMQP specification intends.

Header values

Header values are already readable text when the connector builds them. amqp.header.value.format only decides how they are stored on the record, and that is what decides how a consumer sees them: Kafka’s default SimpleHeaderConverter Base64 encodes any byte array it is given, so a header stored as BYTES arrives Base64 encoded even though the underlying value was plain text.

Value Result

BYTES (default)

Stored with Schema.OPTIONAL_BYTES_SCHEMA. The release 1.0.0 behaviour. Under SimpleHeaderConverter a consumer sees YXBwbGljYXRpb24vanNvbg==; under ByteArrayConverter, used by the examples below, the raw bytes application/json reach the consumer.

STRING

Stored with Schema.OPTIONAL_STRING_SCHEMA. Under SimpleHeaderConverter a consumer sees application/json.

STRING requires a header converter that accepts string values. Use org.apache.kafka.connect.storage.SimpleHeaderConverter, which is the Kafka Connect default. org.apache.kafka.connect.converters.ByteArrayConverter, used in the examples further down, rejects string header values and the task fails on the first message.

Conversion error handling

A message the connector cannot convert never blocks the queue. amqp.error.handling decides what happens to it:

Value Behaviour

DROP

Log a warning and accept the message. It leaves the broker queue and is not written to Kafka, so its content is lost. This is what an unset amqp.error.handling resolves to, unless amqp.unsupported.type.ignore was set to false.

BROKER_REJECT

Log a warning and reject the message, handing it back to the broker. A broker with a dead-letter address configured routes it there with the original bytes intact. Without one, the message is discarded exactly as with DROP.

FAIL

Log a warning, reject the message and fail the task. Nothing is lost, but the connector stops until the cause is resolved.

When amqp.error.handling is left empty, the deprecated amqp.unsupported.type.ignore decides: true means DROP and false means FAIL. An explicit value always wins, so upgrading a connector that relies on the old boolean does not change its behaviour.

The amqp.unsupported.type.log* properties control what the warning contains: whether it is written at all, under which logger name, and whether it includes the full message or only the message id.

Included transformation: DecodeBase64JsonField

Some producers Base64 encode one member of a JSON body a second time, so a CloudEvent arrives with "data":"eyJtcmlkIjoi…​" instead of an object. The connector treats a message body as opaque text and never looks inside it, so this Single Message Transformation restores that member after conversion. It ships inside the connector, so no extra plugin is needed.

transforms=decodeData
transforms.decodeData.type=io.axual.connect.plugins.amqp.transform.DecodeBase64JsonField
transforms.decodeData.field.path=notification.data
Property Description Type Default Importance

field.path

Dot separated path to the Base64 encoded field, for example data or body.payload.data. Required, there is no default. Every segment must be non-empty. Field names that contain a dot cannot be addressed. When a document repeats a key, the first match is decoded.

string

high

With the configuration above:

before   {"notification":{"specversion":"1.0","data":"eyJtcmlkIjoiMDE5ZjhkZjkifQ=="}}
after    {"notification":{"specversion":"1.0","data":{"mrid":"019f8df9"}}}

Only the matched field is rewritten. The rest of the document is spliced back from the original text, so field order, spacing and number formatting are preserved, as is the formatting of the decoded fragment.

Nothing a single record can hold stops the stream. A value that is not a string, a value that is not a JSON object, a field that is not Base64, and content that does not decode to exactly one JSON document all pass through exactly as they arrived, with the reason logged at WARN. A null record value passes through untouched and is not logged.

A record that does not hold the field also passes through, which is normal on a queue where only some records carry it. So does a record whose field is present but null, because an optional member is nothing to decode. Both are logged the same way: the first such record is logged at INFO, once per task, and the rest at DEBUG, so a mistyped field.path does not stay invisible without flooding the log.

Getting Started

This section walks you through configuring the AMQP 1.0 Source Connector on Axual to stream messages from an AMQP broker into a local Kafka stream.

Prerequisites

AMQP broker

  • You have access to an AMQP 1.0 broker (e.g., RabbitMQ or ActiveMQ Artemis) reachable from the Kafka Connect cluster.

  • You have the broker hostname, port, and credentials available.

  • For OAuth2 authentication: you have a valid OAuth2 token endpoint URL, client ID, and client secret (or private key).

Axual stream

The local stream where the connector will produce events must already exist in Axual Self-Service. See Creating topics if you need to create it.

Steps

Step 1 - Create a connector application

  1. Follow the Creating topics documentation in order to create one stream and deploy it onto an environment.
    The name of the stream will be my_amqp_source_topic.
    The key/value types will be String/String.

  2. Follow the Configure and install a connector documentation to set up a new Connector-Application.
    Let’s call it my_amqp_source.
    The plugin name is {plugin-name}.
    If a plugin isn’t available, ask a platform operator to install it on the cluster. A newly installed plugin stays unavailable until someone runs refresh the cluster’s plugin list.

Step 2 - Configure the connector

Always set key.converter, value.converter, and header.converter explicitly. When left unset, Kafka Connect falls back to the worker-level defaults, which vary between installations and may produce unexpected serialisation (e.g., wrapping plain strings in a JSON schema envelope). The AMQP 1.0 connector always emits keys and values as plain strings, so StringConverter is the right choice for both. Headers are raw bytes by default, which ByteArrayConverter handles. If you set amqp.header.value.format=STRING, use org.apache.kafka.connect.storage.SimpleHeaderConverter instead: ByteArrayConverter rejects string header values and the task fails.

  1. Provide the following minimal configuration to connect to your AMQP broker using the Core connector (anonymous or username/password). For OAuth2 authentication, see the OAuth2 configuration section below.

    connector.class

    io.axual.connect.plugins.amqp.source.AmqpSourceConnector

    amqp.host

    Hostname of the AMQP broker
    broker.example.com

    amqp.port

    5672
    Default AMQP port. Use 5671 for TLS.

    amqp.source

    The queue or topic address to consume from
    my.queue

    amqp.auth.username

    Username for SASL PLAIN authentication. Omit for anonymous access.

    amqp.auth.password

    Password for SASL PLAIN authentication. Omit for anonymous access.

    amqp.consumer.qos

    AT_LEAST_ONCE
    Quality of Service level. Options: AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE.

    topic

    my_amqp_source_topic
    The Kafka topic to write messages to.

    amqp.key.extractor

    MESSAGE_ID,SUBJECT,NULL
    Key extraction chain; tries each extractor in order.

    key.converter

    org.apache.kafka.connect.storage.StringConverter

    value.converter

    org.apache.kafka.connect.storage.StringConverter

    header.converter

    org.apache.kafka.connect.converters.ByteArrayConverter

OAuth2 configuration

When using the OAuth2 connector, replace connector.class with io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector and add the following properties:

amqp.oauth.username

The AMQP username to authenticate with. The token is sent as the password.

amqp.oauth.token.endpoint

Access token endpoint URL
https://auth.example.com/oauth2/token

amqp.oauth.method

client_secret_basic
One of: client_secret_basic, client_secret_post, private_key_jwt.

amqp.oauth.client.id

Your OAuth2 client ID

amqp.oauth.client.secret

Your OAuth2 client secret (for client_secret* methods)_

amqp.oauth.scope

Optional space-separated list of requested scopes
read write

  1. Authorize the my_amqp_source source Connector-Application to produce to the my_amqp_source_topic stream.

Step 3 - Start the connector

Start the connector application from Axual Self-Service.

Step 4 - Verify

In Axual Self-Service, use stream-browse on my_amqp_source_topic to confirm messages from the AMQP broker are arriving.

Cleanup

When you are done:

  1. Stop the connector application in Axual Self-Service.

  2. Remove stream access for the application if it is no longer needed.

Examples

The following end-to-end examples show complete configurations for common authentication modes.

Anonymous (no authentication)

{
  "name": "my-amqp-source",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.source.AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue1",
    "amqp.consumer.qos": "AT_LEAST_ONCE",
    "amqp.consumer.durable": "true",
    "topic": "amqp.queue1",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "amqp.validate.connection": "false"
  }
}

Username and password (PLAIN)

{
  "name": "my-amqp-source-plain",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.source.AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue2",
    "amqp.auth.username": "amqp_user",
    "amqp.auth.password": "Amqp@Password2024",
    "amqp.consumer.qos": "AT_LEAST_ONCE",
    "amqp.consumer.durable": "true",
    "topic": "amqp.queue2",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,UUID",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "amqp.validate.connection": "false"
  }
}

JSON body with readable headers

The producer puts the JSON in an AMQP Data section, and the topic should hold that JSON as it was sent, with readable headers. Note the header converter: STRING headers need SimpleHeaderConverter.

{
  "name": "my-amqp-source-json",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.source.AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue3",
    "amqp.consumer.qos": "AT_LEAST_ONCE",
    "amqp.consumer.durable": "true",
    "topic": "amqp.queue3",
    "amqp.data.mode": "JSON",
    "amqp.header.value.format": "STRING",
    "amqp.error.handling": "DROP",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.storage.SimpleHeaderConverter"
  }
}

Add the bundled transformation when the producer also encodes one field inside that JSON:

{
    "transforms": "decodeData",
    "transforms.decodeData.type": "io.axual.connect.plugins.amqp.transform.DecodeBase64JsonField",
    "transforms.decodeData.field.path": "notification.data"
}

OAuth2 - Client Secret Basic

{
  "name": "my-amqp-source-oauth2-basic",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue-basic",
    "amqp.oauth.username": "amqp_user",
    "amqp.oauth.token.endpoint": "https://auth.example.com/oauth2/token",
    "amqp.oauth.method": "client_secret_basic",
    "amqp.oauth.client.id": "my-client",
    "amqp.oauth.client.secret": "Cl13nt$ecret2024",
    "amqp.oauth.scope": "read write",
    "topic": "amqp.oauth.basic",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"
  }
}

OAuth2 - Client Secret Post

{
  "name": "my-amqp-source-oauth2-post",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue-post",
    "amqp.oauth.username": "amqp_user",
    "amqp.oauth.token.endpoint": "https://auth.example.com/oauth2/token",
    "amqp.oauth.method": "client_secret_post",
    "amqp.oauth.client.id": "my-client",
    "amqp.oauth.client.secret": "Cl13nt$ecret2024",
    "amqp.oauth.scope": "read write",
    "topic": "amqp.oauth.post",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"
  }
}

OAuth2 - Private Key JWT with JWK JSON

{
  "name": "my-amqp-source-oauth2-pkjwt-jwk",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue-pkjwt",
    "amqp.oauth.username": "amqp_user",
    "amqp.oauth.token.endpoint": "https://auth.example.com/oauth2/token",
    "amqp.oauth.method": "private_key_jwt",
    "amqp.oauth.client.id": "my-service",
    "amqp.oauth.private.key.jwk.json": "{\"kty\":\"RSA\",\"d\":\"...\",\"n\":\"...\",\"e\":\"AQAB\",\"kid\":\"my-key\"}",
    "amqp.oauth.private.key.jws.algorithm": "RS256",
    "amqp.oauth.private.key.claims.issuer": "my-service",
    "topic": "amqp.oauth.pkjwt.jwk",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"
  }
}

OAuth2 - Private Key JWT with Encrypted PEM

{
  "name": "my-amqp-source-oauth2-pkjwt-pem",
  "config": {
    "connector.class": "io.axual.connect.plugins.amqp.oauth2.source.OAuth2AmqpSourceConnector",
    "tasks.max": "1",
    "amqp.host": "broker.example.com",
    "amqp.port": "5672",
    "amqp.source": "queue-pkjwt-pem",
    "amqp.oauth.username": "amqp_user",
    "amqp.oauth.token.endpoint": "https://auth.example.com/oauth2/token",
    "amqp.oauth.method": "private_key_jwt",
    "amqp.oauth.client.id": "my-service",
    "amqp.oauth.private.key.pem.content": "-----BEGIN ENCRYPTED PRIVATE KEY-----\nMIIF...\n-----END ENCRYPTED PRIVATE KEY-----",
    "amqp.oauth.private.key.pem.password": "changeit",
    "amqp.oauth.private.key.jws.algorithm": "RS256",
    "amqp.oauth.private.key.claims.issuer": "my-service",
    "topic": "amqp.oauth.pkjwt.pem",
    "amqp.key.extractor": "MESSAGE_ID,SUBJECT,NULL",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "header.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"
  }
}

Known limitations

  • AMQP AmqpSequence message body sections are not supported - messages using them result in a conversion error. Data and Value sections are both supported.

  • Array, List and Map values are not supported as a message body or as a record key - they result in a conversion error. As application property values they are silently skipped.

  • The connector always emits string keys and values - schema-aware converters (e.g., Avro, Protobuf) cannot be used.

  • EXACTLY_ONCE QoS requires broker-side support and is not guaranteed on all brokers.

  • BROKER_REJECT error handling only preserves a rejected message if the broker has a dead-letter address configured for that queue. Without one it behaves like DROP.

License

This connector is licensed under a proprietary client-only license. It is provided exclusively to authorized clients and may not be redistributed, modified, or used outside of the terms agreed with Axual.