AMQP 1.0 Source Connectors
Type |
Source |
Core Class |
|
OAuth2 Class |
|
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 ( |
Username and password (SASL PLAIN) |
Core connector. Set |
OAuth2 (client credentials or private key JWT) |
OAuth2 connector ( |
OAuth2 method selection (amqp.oauth.method):
| Method | Description |
|---|---|
|
Client ID and secret sent via HTTP Basic. |
|
Client ID and secret sent in the token request body. |
|
Signed client assertion using a private key. |
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, orNULL -
Message bodies from an AMQP
Datasection orValuesection, 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 |
|---|---|---|---|---|---|
|
The hostname where the AMQP server can be reached. |
string |
non-empty string |
high |
|
|
The port number where the AMQP server listens. |
int |
|
[1,…] |
medium |
|
Enable TLS connection to the AMQP server. |
boolean |
|
medium |
|
|
When set to |
boolean |
|
high |
|
|
Username for connecting to the AMQP server. Uses anonymous authentication if not set. |
string |
null |
high |
|
|
Password for connecting to the AMQP server. |
password |
null |
high |
|
|
The AMQP queue or topic address to consume messages from. This corresponds to the queue name in the broker. |
string |
non-empty string |
high |
|
|
Quality of Service level: |
string |
|
AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE |
medium |
|
Use a durable consumer. The AMQP link remains in place after disconnection so messages are not lost. Essential for persistent topic subscriptions. |
boolean |
|
medium |
|
|
Explicit link name. Overrides auto-generation. When multiple tasks share the same link name, messages are load-balanced between them. |
string |
null |
low |
|
|
Prefix for auto-generated link names. The actual link name will be |
string |
null |
medium |
|
|
Number of messages prefetched and buffered locally before requiring acknowledgments. |
int |
|
[1,…] |
medium |
|
Maximum time in milliseconds to wait for a message when polling. |
long |
|
[100,…] |
low |
|
Maximum time in milliseconds to wait for graceful consumer shutdown. |
long |
|
[1000,…] |
low |
|
Session-level incoming window size for flow control. |
int |
|
[1,…] |
low |
|
Session-level outgoing window size for flow control. |
int |
|
[1,…] |
low |
|
Automatically acknowledge messages after successful processing. |
boolean |
|
medium |
|
|
On processing error: if |
boolean |
|
medium |
|
|
Maximum idle time in milliseconds before the connection is considered dead. |
long |
|
[1,…] |
medium |
|
Maximum size in bytes for individual AMQP frames. |
long |
|
[1,…] |
medium |
|
Network buffer size in bytes for receiving frames. |
int |
|
[1,…] |
medium |
|
Extension size in bytes when the input buffer is insufficient to receive a frame. |
int |
|
[1,…] |
medium |
|
Network buffer size in bytes for sending frames. |
int |
|
[1,…] |
medium |
|
Extension size in bytes when the output buffer is insufficient to send a frame. |
int |
|
[1,…] |
medium |
|
Prefix for auto-generated AMQP container IDs. The actual ID will be |
string |
null |
medium |
|
|
Explicit container ID. Overrides auto-generation. |
string |
null |
low |
|
|
Ordered list of key extraction strategies. Each is tried in turn until a key is found. |
list |
|
high |
|
|
How an encoded message body becomes the Kafka record value. See Message body handling. |
string |
|
BASE64, JSON, BASE64_JSON |
high |
|
How the |
string |
|
BYTES, STRING |
medium |
|
What happens to a message that could not be converted. See Conversion error handling. |
string |
DROP, BROKER_REJECT, FAIL |
medium |
|
|
Deprecated, superseded by |
boolean |
|
medium |
|
|
Log a message that could not be converted. |
boolean |
|
medium |
|
|
Logger name to use when logging a message that could not be converted. Empty means the connector’s own logger. |
string |
low |
||
|
Include the full AMQP message in that log entry. |
boolean |
|
low |
|
|
Include the message id in that log entry. Ignored when the full message is already included. |
boolean |
|
low |
|
|
The Kafka topic to write messages to. |
string |
non-empty string |
high |
|
|
How long to wait for a converted message to be confirmed written to Kafka. While messages are waiting, the AMQP consumer is paused. |
long |
|
[1,…] |
medium |
|
Attempt a live broker connection during connector validation. |
boolean |
|
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 |
|---|---|---|---|---|
|
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 |
|
|
The URL for the access token endpoint. |
string |
high |
|
|
Client authentication method: |
string |
|
medium |
|
The client identifier registered with the authorization server. |
string |
high |
|
|
The client secret (for |
password |
null |
medium |
|
Space-separated list of requested permissions (e.g., |
string |
null |
medium |
|
Private key in JWK JSON format (for |
password |
null |
medium |
|
Private key in PEM encoding (alternative to JWK). |
password |
null |
medium |
|
Password for an encrypted PEM private key. |
password |
null |
low |
|
JWS algorithm used to sign the client assertion. Default is |
string |
|
medium |
|
JWT Issuer Claim ( |
string |
null |
medium |
|
JWT Audience Claim ( |
string |
null |
medium |
|
JWT Expiration Time Claim ( |
long |
|
medium |
|
JWT Not Before Claim ( |
long |
null |
medium |
|
When |
boolean |
|
medium |
|
When |
boolean |
|
medium |
Data conversion
The connector converts AMQP 1.0 messages to Kafka Connect records as follows:
-
Message body (payload) - always emitted as a
Stringvalue usingSchema.OPTIONAL_STRING_SCHEMA.AMQPNullproduces a null record value. How an encoded body becomes that string is set byamqp.data.mode, see Message body handling. -
Record key - extracted by the
amqp.key.extractorchain 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.typeheader, with the valueDataorValue.
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 |
|---|---|
|
|
|
|
|
Decimal notation, as |
|
|
|
Epoch milliseconds (long as string) |
|
|
|
Single-character Unicode string |
|
Canonical UUID text, as |
|
The string itself |
Numeric types ( |
|
|
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 |
|---|---|---|---|
|
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 |
|
|
Bytes decoded as UTF-8 and published unchanged |
Read as Base64 first, then decoded as UTF-8 and published unchanged |
|
Published unchanged |
Published unchanged, text needs no decoding |
Read as Base64, decoded and published unchanged |
|
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 |
|---|---|
|
Stored with |
|
Stored with |
|
|
Conversion error handling
A message the connector cannot convert never blocks the queue. amqp.error.handling decides what
happens to it:
| Value | Behaviour |
|---|---|
|
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 |
|
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 |
|
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 |
|---|---|---|---|---|
|
Dot separated path to the Base64 encoded field, for example |
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
-
Follow the Creating topics documentation in order to create one stream and deploy it onto an environment.
The name of the stream will bemy_amqp_source_topic.
The key/value types will beString/String. -
Follow the Configure and install a connector documentation to set up a new Connector-Application.
Let’s call itmy_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 |
-
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.classio.axual.connect.plugins.amqp.source.AmqpSourceConnectoramqp.hostHostname of the AMQP broker
broker.example.comamqp.port5672
Default AMQP port. Use5671for TLS.amqp.sourceThe queue or topic address to consume from
my.queueamqp.auth.usernameUsername for SASL PLAIN authentication. Omit for anonymous access.
amqp.auth.passwordPassword for SASL PLAIN authentication. Omit for anonymous access.
amqp.consumer.qosAT_LEAST_ONCE
Quality of Service level. Options:AT_MOST_ONCE,AT_LEAST_ONCE,EXACTLY_ONCE.topicmy_amqp_source_topic
The Kafka topic to write messages to.amqp.key.extractorMESSAGE_ID,SUBJECT,NULL
Key extraction chain; tries each extractor in order.key.converterorg.apache.kafka.connect.storage.StringConvertervalue.converterorg.apache.kafka.connect.storage.StringConverterheader.converterorg.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:
|
The AMQP username to authenticate with. The token is sent as the password. |
|
Access token endpoint URL |
|
|
|
Your OAuth2 client ID |
|
Your OAuth2 client secret (for |
|
Optional space-separated list of requested scopes |
-
Authorize the
my_amqp_sourcesource Connector-Application to produce to themy_amqp_source_topicstream.
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
AmqpSequencemessage body sections are not supported - messages using them result in a conversion error.DataandValuesections 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_ONCEQoS requires broker-side support and is not guaranteed on all brokers. -
BROKER_REJECTerror handling only preserves a rejected message if the broker has a dead-letter address configured for that queue. Without one it behaves likeDROP.