Camel Spring Boot

Kafka Share

Consume messages from Apache Kafka topics as a queue, using a share group.

What’s inside

Please refer to the above links for usage and configuration details.

Maven coordinates

<dependency>
    <groupId>org.apache.camel.springboot</groupId>
    <artifactId>camel-kafka-share-starter</artifactId>
</dependency>

Spring Boot Auto-Configuration

The starter supports 80 options, which are listed below.

Name Description Default Type

camel.component.kafka-share.acquire-mode

How the share consumer acquires records. With record_limit, a poll() returns at most maxPollRecords records. With batch_optimized, a poll() can return more records than maxPollRecords, to align with the batches of the topic.

batch_optimized

String

camel.component.kafka-share.additional-properties

Sets additional properties for either kafka consumer or kafka producer in case they can’t be set directly on the camel configurations (e.g.: new Kafka properties that are not reflected yet in Camel configurations), the properties have to be prefixed with additionalProperties., e.g.: additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro. If the properties are set in the application.properties file, they must be prefixed with camel.component.kafka.additional-properties (camel.component.kafka-share.additional-properties for the Kafka Share component) followed by the property name enclosed in square brackets, for example the delivery.timeout.ms property in square brackets. This is a multi-value option with prefix: additionalProperties.

Object>

camel.component.kafka-share.autowired-enabled

Whether autowiring is enabled. This is used for automatic autowiring options (the option must be marked as autowired) by looking up in the registry to find if there is a single instance of matching type, which then gets configured on the component. This can be used for automatic configuring JDBC data sources, JMS connection factories, AWS Clients, etc.

true

Boolean

camel.component.kafka-share.bridge-error-handler

Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming messages, or the likes, will now be processed as a message and handled by the routing Error Handler. Important: This is only possible if the 3rd party component allows Camel to be alerted if an exception was thrown. Some components handle this internally only, and therefore bridgeErrorHandler is not possible. In other situations we may improve the Camel component to hook into the 3rd party component and make this possible for future releases. By default the consumer will use the org.apache.camel.spi.ExceptionHandler to deal with exceptions, that will be logged at WARN or ERROR level and ignored.

false

Boolean

camel.component.kafka-share.brokers

URL of the Kafka brokers to use. The format is host1:port1,host2:port2, and the list can be a subset of brokers or a VIP pointing to a subset of brokers. This option is known as bootstrap.servers in the Kafka documentation.

String

camel.component.kafka-share.check-crcs

Automatically check the CRC32 of the records consumed. This ensures no on-the-wire or on-disk corruption to the messages occurred.

true

Boolean

camel.component.kafka-share.client-id

The client id is a user-specified string sent in each request to help trace calls. It should logically identify the application making the request.

String

camel.component.kafka-share.commit-mode

How the acknowledgements of a poll are committed to the broker. SYNC commits them after the records of the poll are processed, and waits for the result, so a failure to commit is reported to the exception handler. ASYNC commits them without waiting, and a failure to commit is logged.

sync

KafkaShareCommitMode

camel.component.kafka-share.commit-timeout-ms

The maximum time to wait for the acknowledgements to be committed, when commitMode is SYNC. The option is a java.lang.Long type.

5000

Long

camel.component.kafka-share.configuration

Allows to pre-configure the Kafka share component with common options that the endpoints will reuse. The option is a org.apache.camel.component.kafka.share.KafkaShareConfiguration type.

KafkaShareConfiguration

camel.component.kafka-share.connection-max-idle-ms

Close idle connections after the number of milliseconds specified by this config.

540000

Integer

camel.component.kafka-share.consumer-request-timeout-ms

The configuration controls the maximum amount of time the client will wait for the response of a request. If the response is not received before the timeout elapses, the client will resend the request if necessary or fail the request if retries are exhausted.

30000

Integer

camel.component.kafka-share.consumers-count

The number of consumers that connect to the Kafka server. Each consumer runs on its own thread and receives records, as the records of a partition are shared by all the consumers of a share group. Unlike a consumer group, the number of consumers is not limited by the number of partitions.

1

Integer

camel.component.kafka-share.create-consumer-backoff-interval

The delay in millis seconds to wait before trying again to create the kafka consumer (kafka-client).

5000

Long

camel.component.kafka-share.create-consumer-backoff-max-attempts

Maximum attempts to create the kafka consumer (kafka-client), before eventually giving up and failing. Error during creating the consumer may be fatal due to invalid configuration and as such recovery is not possible. However, one part of the validation is DNS resolution of the bootstrap broker hostnames. This may be a temporary networking problem, and could potentially be recoverable. While other errors are fatal, such as some invalid kafka configurations. Unfortunately, kafka-client does not separate this kind of errors. Camel will by default retry forever, and therefore never give up. If you want to give up after many attempts then set this option and Camel will then when giving up terminate the consumer. To try again, you can manually restart the consumer by stopping, and starting the route.

Integer

camel.component.kafka-share.enabled

Whether to enable auto configuration of the kafka-share component. This is enabled by default.

Boolean

camel.component.kafka-share.fetch-max-bytes

The maximum amount of data the server should return for a fetch request. This is not an absolute maximum: if the first record batch in the first non-empty partition of the fetch is larger than this value, the record batch will still be returned to ensure that the consumer can make progress.

52428800

Integer

camel.component.kafka-share.fetch-min-bytes

The minimum amount of data the server should return for a fetch request. If insufficient data is available, the request will wait for that much data to accumulate before answering the request.

1

Integer

camel.component.kafka-share.fetch-wait-max-ms

The maximum amount of time the server will block before answering the fetch request if there isn’t enough data to immediately satisfy fetch.min.bytes.

500

Integer

camel.component.kafka-share.group-id

The name of the share group. All the consumers that use the same share group name share the records of the topics: each record is delivered to one of them.

String

camel.component.kafka-share.header-deserializer

To use a custom KafkaHeaderDeserializer to deserialize kafka headers values. The option is a org.apache.camel.component.kafka.serde.KafkaHeaderDeserializer type.

KafkaHeaderDeserializer

camel.component.kafka-share.header-filter-strategy

To use a custom HeaderFilterStrategy to filter header to and from Camel message. The option is a org.apache.camel.spi.HeaderFilterStrategy type.

HeaderFilterStrategy

camel.component.kafka-share.health-check-consumer-enabled

Used for enabling or disabling all consumer based health checks from this component

true

Boolean

camel.component.kafka-share.health-check-producer-enabled

Used for enabling or disabling all producer based health checks from this component. Notice: Camel has by default disabled all producer based health-checks. You can turn on producer checks globally by setting camel.health.producersEnabled=true.

true

Boolean

camel.component.kafka-share.kafka-share-client-factory

Factory to use for creating org.apache.kafka.clients.consumer.KafkaShareConsumer instances. This allows configuring a custom factory to create instances with logic that extends the vanilla Kafka clients. The option is a org.apache.camel.component.kafka.share.KafkaShareClientFactory type.

KafkaShareClientFactory

camel.component.kafka-share.kerberos-before-relogin-min-time

Login thread sleep time between refresh attempts.

60000

Integer

camel.component.kafka-share.kerberos-config-location

Location of the kerberos config file.

String

camel.component.kafka-share.kerberos-init-cmd

Kerberos kinit command path. Default is /usr/bin/kinit

/usr/bin/kinit

String

camel.component.kafka-share.kerberos-principal-to-local-rules

A list of rules for mapping from principal names to short names (typically operating system usernames). The rules are evaluated in order, and the first rule that matches a principal name is used to map it to a short name. Any later rules in the list are ignored. By default, principal names of the form {username}/{hostname}{REALM} are mapped to {username}. For more details on the format, please see the Security Authorization and ACLs documentation (at the Apache Kafka project website). Multiple values can be separated by comma

DEFAULT

String

camel.component.kafka-share.kerberos-renew-jitter

Percentage of random jitter added to the renewal time.

Double

camel.component.kafka-share.kerberos-renew-window-factor

Login thread will sleep until the specified window factor of time from last refresh to ticket’s expiry has been reached, at which time it will try to renew the ticket.

Double

camel.component.kafka-share.key-deserializer

Deserializer class for the key that implements the Deserializer interface.

org.apache.kafka.common.serialization.StringDeserializer

String

camel.component.kafka-share.max-partition-fetch-bytes

The maximum amount of data per-partition the server will return.

1048576

Integer

camel.component.kafka-share.max-poll-records

The maximum number of records returned in a single call to poll(). With the batch_optimized acquire mode, a poll can return more records, to align with the batches of the topic.

500

Integer

camel.component.kafka-share.metadata-max-age-ms

The period of time in milliseconds after which we force a refresh of metadata even if we haven’t seen any partition leadership changes to proactively discover any new brokers or partitions.

300000

Integer

camel.component.kafka-share.metric-reporters

A list of classes to use as metrics reporters. Implementing the MetricReporter interface allows plugging in classes that will be notified of new metric creation. The JmxReporter is always included to register JMX statistics.

String

camel.component.kafka-share.metrics-sample-window-ms

The window of time a metrics sample is computed over.

30000

Integer

camel.component.kafka-share.no-of-metrics-sample

The number of samples maintained to compute metrics.

2

Integer

camel.component.kafka-share.oauth-client-id

OAuth client ID. Used when saslAuthType is set to OAUTH.

String

camel.component.kafka-share.oauth-client-secret

OAuth client secret. Used when saslAuthType is set to OAUTH.

String

camel.component.kafka-share.oauth-scope

OAuth scope. Used when saslAuthType is set to OAUTH.

String

camel.component.kafka-share.oauth-token-endpoint-uri

OAuth token endpoint URI. Used when saslAuthType is set to OAUTH.

String

camel.component.kafka-share.on-failure

How to acknowledge a record whose exchange failed or was rolled back. RELEASE makes the record available again, to this or another consumer, until the broker delivery count limit is reached. REJECT discards the record. ACCEPT marks the record as consumed.

release

KafkaShareAcknowledgeType

camel.component.kafka-share.poll-exception-strategy

To use a custom strategy with the consumer to control how to handle exceptions thrown from the Kafka broker while polling messages. The option is a org.apache.camel.component.kafka.PollExceptionStrategy type.

PollExceptionStrategy

camel.component.kafka-share.poll-on-error

What to do if the share consumer throws an exception while polling for new records. DISCARD and RETRY log the exception and poll again: unlike the kafka component, there is no record to skip or to retry, as the poll itself failed, and the records that fail in the route are handled with onFailure. ERROR_HANDLER lets the exception handler of the consumer handle the exception, and polls again. RECONNECT closes the share consumer and creates a new one. STOP stops consuming. An authentication or authorization failure always stops consuming.

error-handler

PollOnError

camel.component.kafka-share.poll-timeout-ms

The timeout used when polling the share consumer. The option is a java.lang.Long type.

5000

Long

camel.component.kafka-share.pre-validate-host-and-port

Whether to eager validate that broker host:port is valid and can be DNS resolved to known host during starting this consumer. If the validation fails, then an exception is thrown, which makes Camel fail fast. Disabling this will postpone the validation after the consumer is started, and Camel will keep re-connecting in case of validation or DNS resolution error.

true

Boolean

camel.component.kafka-share.receive-buffer-bytes

The size of the TCP receive buffer (SO_RCVBUF) to use when reading data.

65536

Integer

camel.component.kafka-share.reconnect-backoff-max-ms

The maximum amount of time in milliseconds to wait when reconnecting to a broker that has repeatedly failed to connect. If provided, the backoff per host will increase exponentially for each consecutive connection failure, up to this maximum. After calculating the backoff increase, 20% random jitter is added to avoid connection storms.

1000

Integer

camel.component.kafka-share.reconnect-backoff-ms

The amount of time to wait before attempting to reconnect to a given host. This avoids repeatedly connecting to a host in a tight loop. This backoff applies to all requests sent by the consumer to the broker.

50

Integer

camel.component.kafka-share.retry-backoff-max-ms

The maximum amount of time in milliseconds to wait when retrying a request to the broker that has repeatedly failed. If provided, the backoff per client will increase exponentially for each failed request, up to this maximum. To prevent all clients from being synchronized upon retry, a randomized jitter with a factor of 0.2 will be applied to the backoff, resulting in the backoff falling within a range between 20% below and 20% above the computed value. If retry.backoff.ms is set to be higher than retry.backoff.max.ms, then retry.backoff.max.ms will be used as a constant backoff from the beginning without any exponential increase

1000

Integer

camel.component.kafka-share.retry-backoff-ms

The amount of time to wait before attempting to retry a failed request to a given topic partition. This avoids repeatedly sending requests in a tight loop under some failure scenarios. This value is the initial backoff value and will increase exponentially for each failed request, up to the retry.backoff.max.ms value.

100

Integer

camel.component.kafka-share.sasl-auth-type

Simplified authentication type to use. This provides an easier way to configure Kafka authentication without manually setting securityProtocol, saslMechanism, and saslJaasConfig. When set, the appropriate security settings are automatically derived. Note: This is optional. You can still use the traditional approach with explicit securityProtocol, saslMechanism, and saslJaasConfig properties.

KafkaAuthType

camel.component.kafka-share.sasl-jaas-config

Expose the kafka sasl.jaas.config parameter Example: org.apache.kafka.common.security.plain.PlainLoginModule required username=USERNAME password=PASSWORD;

String

camel.component.kafka-share.sasl-kerberos-service-name

The Kerberos principal name that Kafka runs as. This can be defined either in Kafka’s JAAS config or in Kafka’s config.

String

camel.component.kafka-share.sasl-mechanism

The Simple Authentication and Security Layer (SASL) Mechanism used. For the valid values see http://www.iana.org/assignments/sasl-mechanisms/sasl-mechanisms.xhtml

GSSAPI

String

camel.component.kafka-share.sasl-password

Password for SASL authentication. Used when saslAuthType is set to PLAIN, SCRAM_SHA_256, or SCRAM_SHA_512.

String

camel.component.kafka-share.sasl-username

Username for SASL authentication. Used when saslAuthType is set to PLAIN, SCRAM_SHA_256, or SCRAM_SHA_512.

String

camel.component.kafka-share.schema-registry-u-r-l

URL of the schema registry servers to use. The format is host1:port1,host2:port2. This is known as schema.registry.url in multiple Schema registries documentation. This option is only available externally (not standard Apache Kafka)

String

camel.component.kafka-share.security-protocol

Protocol used to communicate with brokers. SASL_PLAINTEXT, PLAINTEXT, SASL_SSL and SSL are supported

PLAINTEXT

String

camel.component.kafka-share.send-buffer-bytes

Socket write buffer size

131072

Integer

camel.component.kafka-share.shutdown-timeout

Timeout in milliseconds to wait gracefully for the consumer or producer to shut down and terminate its worker threads.

30000

Integer

camel.component.kafka-share.specific-avro-reader

This enables the use of a specific Avro reader for use with the in multiple Schema registries documentation with Avro Deserializers implementation. This option is only available externally (not standard Apache Kafka)

false

Boolean

camel.component.kafka-share.ssl-cipher-suites

A list of cipher suites. This is a named combination of authentication, encryption, MAC and key exchange algorithm used to negotiate the security settings for a network connection using TLS or SSL network protocol. By default, all the available cipher suites are supported.

String

camel.component.kafka-share.ssl-context-parameters

SSL configuration using a Camel SSLContextParameters object. If configured, it’s applied before the other SSL endpoint parameters. NOTE: Kafka only supports loading keystore from file locations, so prefix the location with file: in the KeyStoreParameters.resource option. The option is a org.apache.camel.support.jsse.SSLContextParameters type.

SSLContextParameters

camel.component.kafka-share.ssl-enabled-protocols

The list of protocols enabled for SSL connections. The default is TLSv1.2,TLSv1.3 when running with Java 11 or newer, TLSv1.2 otherwise. With the default value for Java 11, clients and servers will prefer TLSv1.3 if both support it and fallback to TLSv1.2 otherwise (assuming both support at least TLSv1.2). This default should be fine for most cases. Also see the config documentation for SslProtocol.

String

camel.component.kafka-share.ssl-endpoint-algorithm

The endpoint identification algorithm to validate server hostname using server certificate. Use none or false to disable server hostname verification.

https

String

camel.component.kafka-share.ssl-key-password

The password of the private key in the key store file or the PEM key specified in sslKeystoreKey. This is required for clients only if two-way authentication is configured.

String

camel.component.kafka-share.ssl-keymanager-algorithm

The algorithm used by key manager factory for SSL connections. Default value is the key manager factory algorithm configured for the Java Virtual Machine.

SunX509

String

camel.component.kafka-share.ssl-keystore-location

The location of the key store file. This is optional for the client and can be used for two-way authentication for the client.

String

camel.component.kafka-share.ssl-keystore-password

The store password for the key store file. This is optional for the client and only needed if sslKeystoreLocation is configured. Key store password is not supported for PEM format.

String

camel.component.kafka-share.ssl-keystore-type

The file format of the key store file. This is optional for the client. The default value is JKS

JKS

String

camel.component.kafka-share.ssl-protocol

The SSL protocol used to generate the SSLContext. The default is TLSv1.3 when running with Java 11 or newer, TLSv1.2 otherwise. This value should be fine for most use cases. Allowed values in recent JVMs are TLSv1.2 and TLSv1.3. TLS, TLSv1.1, SSL, SSLv2 and SSLv3 may be supported in older JVMs, but their usage is discouraged due to known security vulnerabilities. With the default value for this config and sslEnabledProtocols, clients will downgrade to TLSv1.2 if the server does not support TLSv1.3. If this config is set to TLSv1.2, clients will not use TLSv1.3 even if it is one of the values in sslEnabledProtocols and the server only supports TLSv1.3.

String

camel.component.kafka-share.ssl-provider

The name of the security provider used for SSL connections. Default value is the default security provider of the JVM.

String

camel.component.kafka-share.ssl-trustmanager-algorithm

The algorithm used by trust manager factory for SSL connections. Default value is the trust manager factory algorithm configured for the Java Virtual Machine.

PKIX

String

camel.component.kafka-share.ssl-truststore-location

The location of the trust store file.

String

camel.component.kafka-share.ssl-truststore-password

The password for the trust store file. If a password is not set, trust store file configured will still be used, but integrity checking is disabled. Trust store password is not supported for PEM format.

String

camel.component.kafka-share.ssl-truststore-type

The file format of the trust store file. The default value is JKS.

JKS

String

camel.component.kafka-share.use-global-ssl-context-parameters

Enable usage of global SSL context parameters.

false

Boolean

camel.component.kafka-share.value-deserializer

Deserializer class for value that implements the Deserializer interface.

org.apache.kafka.common.serialization.StringDeserializer

String